@cryptotaxi247 / netdata-1 / commits / 78c3f35af

Improved ACLK (#8498)

Improved the stability of the ACLK

Stelios Fragkakis committed Mar 26, 2020 at 18:02 UTC 78c3f35af87b57a2988f3878302284a9de37977e
6 files changed +82 -51
aclk/agent_cloud_link.c
+74 -45
@@ -117,7 +117,12 @@ int cloud_to_agent_parse(JSON_ENTRY *e)
117 break;
118 }
119 if (!strcmp(e->name, "payload")) {
120 - data->payload = strdupz(e->data.string);
120 + if (likely(e->data.string)) {
121 + size_t len = strlen(e->data.string);
122 + data->payload = mallocz(len+1);
123 + if (!url_decode_r(data->payload, e->data.string, len + 1))
124 + strcpy(data->payload, e->data.string);
125 + }
126 break;
127 }
128 break;
@@ -302,20 +307,19 @@ int aclk_queue_query(char *topic, char *data, char *msg_id, char *query, int run
307
308 // Ignore all commands while we wait for the agent to initialize
309 if (unlikely(waiting_init))
305 - return 0;
310 + return 1;
311
312 run_after = now_realtime_sec() + run_after;
313
314 QUERY_LOCK;
315 struct aclk_query *last_query = NULL;
316
312 - //last_query = NULL;
317 tmp_query = aclk_query_find(topic, data, msg_id, query, aclk_cmd, &last_query);
318 if (unlikely(tmp_query)) {
319 if (tmp_query->run_after == run_after) {
320 QUERY_UNLOCK;
321 QUERY_THREAD_WAKEUP;
318 - return 0;
322 + return 1;
323 }
324
325 if (last_query)
@@ -509,7 +513,7 @@ static void _free_collector(struct _collector *collector)
513 *
514 */
515 #ifdef ACLK_DEBUG
512 -static void _dump_connector_list()
516 +static void _dump_collector_list()
517 {
518 struct _collector *tmp_collector;
519
@@ -543,7 +547,7 @@ static void _dump_connector_list()
547 * This will cleanup the collector list
548 *
549 */
546 -static void _reset_connector_list()
550 +static void _reset_collector_list()
551 {
552 struct _collector *tmp_collector, *next_collector;
553
@@ -689,8 +693,10 @@ void aclk_add_collector(const char *hostname, const char *plugin_name, const cha
693
694 if (unlikely(agent_state == AGENT_INITIALIZING))
695 last_init_sequence = now_realtime_sec();
692 - else
693 - aclk_queue_query("connector", NULL, NULL, NULL, 0, 1, ACLK_CMD_ONCONNECT);
696 + else {
697 + if (unlikely(aclk_queue_query("collector", NULL, NULL, NULL, 0, 1, ACLK_CMD_ONCONNECT)))
698 + debug(D_ACLK, "ACLK failed to queue on_connect command on collector addition");
699 + }
700
701 COLLECTOR_UNLOCK;
702 }
@@ -724,8 +730,10 @@ void aclk_del_collector(const char *hostname, const char *plugin_name, const cha
730
731 if (unlikely(agent_state == AGENT_INITIALIZING))
732 last_init_sequence = now_realtime_sec();
727 - else
728 - aclk_queue_query("on_connect", NULL, NULL, NULL, 0, 1, ACLK_CMD_ONCONNECT);
733 + else {
734 + if (unlikely(aclk_queue_query("collector", NULL, NULL, NULL, 0, 1, ACLK_CMD_ONCONNECT)))
735 + debug(D_ACLK, "ACLK failed to queue on_connect command on collector deletion");
736 + }
737
738 _free_collector(tmp_collector);
739 }
@@ -757,12 +765,10 @@ int aclk_execute_query(struct aclk_query *this_query)
765
766 aclk_create_header(local_buffer, "http", this_query->msg_id);
767
760 - //if (rc != HTTP_RESP_OK || strcmp(mysep ? mysep + 1 : "noop", "badge.svg") == 0)
768 + char *encoded_response = aclk_encode_response(w->response.data);
769 +
770 buffer_sprintf(
762 - local_buffer, "{\n\"code\": %d,\n\"body\": \"%s\"\n}", rc, aclk_encode_response(w->response.data)->buffer);
763 - //else
764 - // buffer_sprintf(local_buffer, "{\n\"code\": %d,\n\"body\": %s\n}", rc,
765 - // aclk_encode_response(w->response.data)->buffer);
771 + local_buffer, "{\n\"code\": %d,\n\"body\": \"%s\"\n}", rc, encoded_response);
772
773 buffer_sprintf(local_buffer, "\n}");
774
@@ -771,6 +777,7 @@ int aclk_execute_query(struct aclk_query *this_query)
777 buffer_free(w->response.data);
778 freez(w);
779 buffer_free(local_buffer);
780 + freez(encoded_response);
781 return 0;
782 }
783 return 1;
@@ -880,7 +887,7 @@ static void aclk_query_thread_cleanup(void *ptr)
887
888 COLLECTOR_LOCK;
889
883 - _reset_connector_list();
890 + _reset_collector_list();
891 freez(collector_list);
892
893 COLLECTOR_UNLOCK;
@@ -908,7 +915,7 @@ void *aclk_query_main_thread(void *ptr)
915 agent_state = AGENT_STABLE;
916 info("AGENT stable, last collector initialization activity was %ld seconds ago", checkpoint);
917 #ifdef ACLK_DEBUG
911 - _dump_connector_list();
918 + _dump_collector_list();
919 #endif
920 break;
921 }
@@ -919,7 +926,11 @@ void *aclk_query_main_thread(void *ptr)
926 while (!netdata_exit) {
927 if (unlikely(!aclk_metadata_submitted)) {
928 aclk_metadata_submitted = ACLK_METADATA_CMD_QUEUED;
922 - aclk_queue_query("on_connect", NULL, NULL, NULL, 0, 1, ACLK_CMD_ONCONNECT);
929 + if (unlikely(aclk_queue_query("on_connect", NULL, NULL, NULL, 0, 1, ACLK_CMD_ONCONNECT))) {
930 + errno = 0;
931 + error("ACLK failed to queue on_connect command");
932 + aclk_metadata_submitted = 0;
933 + }
934 }
935
936 aclk_process_queries();
@@ -1303,13 +1314,11 @@ void *aclk_main(void *ptr)
1314 last_init_sequence = now_realtime_sec();
1315 query_thread = NULL;
1316
1306 -
1317 char *aclk_hostname = NULL; // Initializers are over-written but prevent gcc complaining about clobbering.
1318 char *aclk_port = NULL;
1319 uint32_t port_num = 0;
1320 char *cloud_base_url = config_get(CONFIG_SECTION_CLOUD, "cloud base url", "https://netdata.cloud");
1311 - if( aclk_decode_base_url(cloud_base_url, &aclk_hostname, &aclk_port))
1312 - {
1321 + if (aclk_decode_base_url(cloud_base_url, &aclk_hostname, &aclk_port)) {
1322 error("Configuration error - cannot use agent cloud link");
1323 return NULL;
1324 }
@@ -1376,7 +1385,7 @@ void *aclk_main(void *ptr)
1385
1386 // TODO: Move to on-connect
1387 if (unlikely(!aclk_subscribed)) {
1379 - aclk_subscribed = !aclk_subscribe(ACLK_COMMAND_TOPIC, 2);
1388 + aclk_subscribed = !aclk_subscribe(ACLK_COMMAND_TOPIC, 1);
1389 }
1390
1391 if (unlikely(!query_thread)) {
@@ -1416,7 +1425,7 @@ int aclk_send_message(char *sub_topic, char *message, char *msg_id)
1425
1426 UNUSED(msg_id);
1427
1419 - if(!aclk_connected)
1428 + if (!aclk_connected)
1429 return 0;
1430
1431 if (unlikely(!message))
@@ -1541,20 +1550,20 @@ inline void aclk_create_header(BUFFER *dest, char *type, char *msg_id)
1550 *
1551 */
1552
1544 -BUFFER *aclk_encode_response(BUFFER *contents)
1553 +char *aclk_encode_response(BUFFER *contents)
1554 {
1555 char *tmp_buffer = mallocz(contents->len * 2);
1556 char *src, *dst;
1557 + size_t content_size = contents->len;
1558
1559 src = contents->buffer;
1560 dst = tmp_buffer;
1551 - while (*src) {
1561 + while (content_size > 0) {
1562 switch (*src) {
1563 case '\n':
1554 - *dst++ = '\\';
1555 - *dst++ = 'n';
1564 + case '\t':
1565 break;
1557 - case 0x01 ... 0x09:
1566 + case 0x01 ... 0x08:
1567 case 0x0b ... 0x1F:
1568 *dst++ = '\\';
1569 *dst++ = '0';
@@ -1563,7 +1572,6 @@ BUFFER *aclk_encode_response(BUFFER *contents)
1572 *dst++ = to_hex(*src);
1573 break;
1574 case '\"':
1566 - case '\'':
1575 *dst++ = '\\';
1576 *dst++ = *src;
1577 break;
@@ -1571,19 +1579,18 @@ BUFFER *aclk_encode_response(BUFFER *contents)
1579 *dst++ = *src;
1580 }
1581 src++;
1582 + content_size--;
1583 }
1584 *dst = '\0';
1585
1577 - buffer_flush(contents);
1578 - buffer_sprintf(contents, "%s", tmp_buffer);
1579 -
1580 - freez(tmp_buffer);
1581 - return contents;
1586 + return tmp_buffer;
1587 }
1588
1589 /*
1585 - * This will send the alarms configuration
1586 - * and
1590 + * This will send alarm information which includes
1591 + * configured alarms
1592 + * alarm_log
1593 + * active alarms
1594 */
1595 void aclk_send_alarm_metadata()
1596 {
@@ -1611,12 +1618,16 @@ void aclk_send_alarm_metadata()
1618
1619 buffer_sprintf(local_buffer, "\n}\n}");
1620 aclk_send_message(ACLK_ALARMS_TOPIC, local_buffer->buffer, msg_id);
1614 - debug(D_ACLK, "Metadata %s encoded has %zu bytes", msg_id, local_buffer->len);
1621
1622 freez(msg_id);
1623 buffer_free(local_buffer);
1624 }
1625
1626 +/*
1627 + * This will send the agent metadata
1628 + * /api/v1/info
1629 + * charts
1630 + */
1631 int aclk_send_info_metadata()
1632 {
1633 BUFFER *local_buffer = buffer_create(NETDATA_WEB_RESPONSE_INITIAL_SIZE);
@@ -1638,9 +1649,8 @@ int aclk_send_info_metadata()
1649 debug(D_ACLK, "Metadata %s with chart has %zu bytes", msg_id, local_buffer->len);
1650
1651 aclk_send_message(ACLK_METADATA_TOPIC, local_buffer->buffer, msg_id);
1641 - debug(D_ACLK, "Metadata %s encoded has %zu bytes", msg_id, local_buffer->len);
1642 - freez(msg_id);
1652
1653 + freez(msg_id);
1654 buffer_free(local_buffer);
1655 return 0;
1656 }
@@ -1691,7 +1701,12 @@ void aclk_alarm_reload()
1701 if (unlikely(agent_state == AGENT_INITIALIZING))
1702 return;
1703
1694 - aclk_queue_query("on_connect", NULL, NULL, NULL, 0, 1, ACLK_CMD_ONCONNECT);
1704 + if (unlikely(aclk_queue_query("on_connect", NULL, NULL, NULL, 0, 1, ACLK_CMD_ONCONNECT))) {
1705 + if (likely(aclk_connected)) {
1706 + errno = 0;
1707 + error("ACLK failed to queue on_connect command on alarm reload");
1708 + }
1709 + }
1710 }
1711 //rrd_stats_api_v1_chart(RRDSET *st, BUFFER *buf)
1712
@@ -1742,8 +1757,14 @@ int aclk_update_chart(RRDHOST *host, char *chart_name, ACLK_CMD aclk_cmd)
1757
1758 if (unlikely(agent_state == AGENT_INITIALIZING))
1759 last_init_sequence = now_realtime_sec();
1745 - else
1746 - aclk_queue_query("_chart", host->hostname, NULL, chart_name, 0, 1, aclk_cmd);
1760 + else {
1761 + if (unlikely(aclk_queue_query("_chart", host->hostname, NULL, chart_name, 0, 1, aclk_cmd))) {
1762 + if (likely(aclk_connected)) {
1763 + errno = 0;
1764 + error("ACLK failed to queue chart_update command");
1765 + }
1766 + }
1767 + }
1768 return 0;
1769 #endif
1770 }
@@ -1779,7 +1800,13 @@ int aclk_update_alarm(RRDHOST *host, ALARM_ENTRY *ae)
1800 netdata_rwlock_unlock(&host->health_log.alarm_log_rwlock);
1801
1802 buffer_sprintf(local_buffer, "\n}");
1782 - aclk_queue_query(ACLK_ALARMS_TOPIC, NULL, msg_id, local_buffer->buffer, 0, 1, ACLK_CMD_ALARM);
1803 +
1804 + if (unlikely(aclk_queue_query(ACLK_ALARMS_TOPIC, NULL, msg_id, local_buffer->buffer, 0, 1, ACLK_CMD_ALARM))) {
1805 + if (likely(aclk_connected)) {
1806 + errno = 0;
1807 + error("ACLK failed to queue alarm_command on alarm_update");
1808 + }
1809 + }
1810
1811 freez(msg_id);
1812 buffer_free(local_buffer);
@@ -1796,6 +1823,7 @@ int aclk_handle_cloud_request(char *payload)
1823 .type_id = NULL, .msg_id = NULL, .callback_topic = NULL, .payload = NULL, .version = 0
1824 };
1825
1826 +
1827 if (unlikely(agent_state == AGENT_INITIALIZING)) {
1828 debug(D_ACLK, "Ignoring cloud request; agent not in stable state");
1829 return 0;
@@ -1806,7 +1834,7 @@ int aclk_handle_cloud_request(char *payload)
1834 return 0;
1835 }
1836
1809 - debug(D_ACLK, "ACLK incoming message [%s]", payload);
1837 + debug(D_ACLK, "ACLK incoming message (%s)", payload);
1838
1839 int rc = json_parse(payload, &cloud_to_agent, cloud_to_agent_parse);
1840
@@ -1835,7 +1863,8 @@ int aclk_handle_cloud_request(char *payload)
1863 return 1;
1864 }
1865
1838 - aclk_submit_request(&cloud_to_agent);
1866 + if (unlikely(aclk_submit_request(&cloud_to_agent)))
1867 + debug(D_ACLK, "ACLK failed to queue incoming message (%s)", payload);
1868
1869 // Note: the payload comes from the callback and it will be automatically freed
1870 return 0;
aclk/agent_cloud_link.h
+1 -1
@@ -106,7 +106,7 @@ void aclk_del_collector(const char *hostname, const char *plugin_name, const cha
106 void aclk_alarm_reload();
107 void aclk_send_alarm_metadata();
108 int aclk_execute_query(struct aclk_query *query);
109 -BUFFER *aclk_encode_response(BUFFER *contents);
109 +char *aclk_encode_response(BUFFER *contents);
110 unsigned long int aclk_reconnect_delay(int mode);
111 extern void health_alarm_entry2json_nolock(BUFFER *wb, ALARM_ENTRY *ae, RRDHOST *host);
112 void aclk_single_update_enable();
aclk/mqtt.c
+1 -1
@@ -322,7 +322,7 @@ int _link_send_message(char *topic, unsigned char *message, int *mid)
322 return rc;
323
324 int msg_len = strlen((char*)message);
325 - error("Sending MQTT len=%d starts %02x %02x %02x", msg_len, message[0], message[1], message[2]);
325 + info("Sending MQTT len=%d starts %02x %02x %02x", msg_len, message[0], message[1], message[2]);
326 rc = mosquitto_publish(mosq, mid, topic, msg_len, message, ACLK_QOS, 0);
327
328 // TODO: Add better handling -- error will flood the logfile here
claim/claim.c
+2 -2
@@ -41,8 +41,8 @@ extern struct registry registry;
41 /* rrd_init() must have been called before this function */
42 void claim_agent(char *claiming_arguments)
43 {
44 -#ifndef ENABLE_ACLK
45 - info("The claiming feature is under development and still subject to change before the next release");
44 +#ifndef ENABLE_CLOUD
45 + info("The claiming feature has been disabled");
46 return;
47 #endif
48
daemon/commands.c
+3 -1
@@ -186,8 +186,10 @@ static cmd_status_t cmd_reload_claiming_state_execute(char *args, char **message
186 (void)args;
187 (void)message;
188
189 - info("The claiming feature is still in development and subject to change before the next release");
189 +#ifndef ENABLE_CLOUD
190 + info("The claiming feature has been disabled");
191 return CMD_STATUS_FAILURE;
192 +#endif
193
194 error_log_limit_unlimited();
195 info("COMMAND: Reloading Agent Claiming configuration.");
web/api/web_api_v1.c
+1 -1
@@ -489,7 +489,7 @@ inline int web_client_api_request_v1_data(RRDHOST *host, struct web_client *w, c
489 st->last_accessed_time = now_realtime_sec();
490
491 long long before = (before_str && *before_str)?str2l(before_str):0;
492 - long long after = (after_str && *after_str) ?str2l(after_str):0;
492 + long long after = (after_str && *after_str) ?str2l(after_str):-600;
493 int points = (points_str && *points_str)?str2i(points_str):0;
494 long group_time = (group_time_str && *group_time_str)?str2l(group_time_str):0;
495