Fix cloud connect after claim (#19547)
Stelios Fragkakis committed
Jan 31, 2025 at 18:04 UTC
d35dcc38236f03f8ac846cd294df1807f629af81
5 files changed
+45
-40
src/aclk/aclk.c
+8
-3
@@ -1002,12 +1002,15 @@ void aclk_host_state_update(RRDHOST *host, int cmd, int queryable)
1002
aclk_add_job(query);
1003
}
1004
1005
-void aclk_send_node_instances()
1005
+void aclk_send_node_instances(mqtt_wss_client client)
1006
{
1007
struct node_instance_list *list_head = get_node_list();
1008
struct node_instance_list *list = list_head;
1009
if (unlikely(!list)) {
1010
error_report("Failure to get_node_list from DB!");
1011
+ sleep_usec(USEC_PER_SEC);
1012
+ aclk_query_t query = aclk_query_new(SEND_NODE_INSTANCES);
1013
+ aclk_add_job(query);
1014
return;
1015
}
1016
while (!uuid_is_null(list->host_id)) {
@@ -1045,7 +1048,8 @@ void aclk_send_node_instances()
1048
freez((void*)node_state_update.node_id);
1049
query->data.bin_payload.msg_name = "UpdateNodeInstanceConnection";
1050
query->data.bin_payload.topic = ACLK_TOPICID_NODE_CONN;
1048
- aclk_add_job(query);
1051
+ send_bin_msg(client, query);
1052
+ aclk_query_free(query);
1053
} else {
1054
aclk_query_t create_query;
1055
create_query = aclk_query_new(REGISTER_NODE);
@@ -1067,7 +1071,8 @@ void aclk_send_node_instances()
1071
(char*)node_instance_creation.machine_guid, list->hops);
1072
1073
freez((void *)node_instance_creation.machine_guid);
1070
- aclk_add_job(create_query);
1074
+ send_bin_msg(client, create_query);
1075
+ aclk_query_free(create_query);
1076
}
1077
freez(list->hostname);
1078
src/aclk/aclk.h
+1
-1
@@ -96,7 +96,7 @@ extern struct aclk_shared_state {
96
void aclk_host_state_update(RRDHOST *host, int cmd, int queryable);
97
bool aclk_host_state_update_auto(RRDHOST *host);
98
99
-void aclk_send_node_instances(void);
99
+void aclk_send_node_instances(mqtt_wss_client client);
100
101
void aclk_send_bin_msg(char *msg, size_t msg_len, enum aclk_topics subtopic, const char *msgname);
102
src/database/sqlite/sqlite_aclk.c
+28
-28
@@ -13,7 +13,7 @@ void sanity_check(void) {
13
#include "../aclk_query.h"
14
#include "../aclk_capas.h"
15
16
-static void create_node_instance_result_job(const char *machine_guid, const char *node_id)
16
+static void create_node_instance_result_job(mqtt_wss_client client __maybe_unused, const char *machine_guid, const char *node_id)
17
{
18
nd_uuid_t host_uuid, node_uuid;
19
@@ -32,33 +32,33 @@ static void create_node_instance_result_job(const char *machine_guid, const char
32
netdata_log_error("Cannot find machine_guid provided by CreateNodeInstanceResult");
33
return;
34
}
35
-
35
sql_update_node_id(&host_uuid, &node_uuid);
37
-
38
- aclk_query_t query = aclk_query_new(NODE_STATE_UPDATE);
39
- node_instance_connection_t node_state_update = {
40
- .hops = 1,
41
- .live = 0,
42
- .queryable = 1,
43
- .session_id = aclk_session_newarch,
44
- .node_id = node_id,
45
- .capabilities = NULL};
46
-
47
- node_state_update.live = rrdhost_is_local(host) ? 1 : 0;
48
- node_state_update.hops = rrdhost_ingestion_hops(host);
49
- node_state_update.capabilities = aclk_get_node_instance_capas(host);
36
schedule_node_state_update(host, 5000);
51
-
52
- CLAIM_ID claim_id = claim_id_get();
53
- node_state_update.claim_id = claim_id_is_set(claim_id) ? claim_id.str : NULL;
54
- query->data.bin_payload.payload = generate_node_instance_connection(&query->data.bin_payload.size, &node_state_update);
55
-
56
- freez((void *)node_state_update.capabilities);
57
-
58
- query->data.bin_payload.msg_name = "UpdateNodeInstanceConnection";
59
- query->data.bin_payload.topic = ACLK_TOPICID_NODE_CONN;
60
-
61
- aclk_add_job(query);
37
+//
38
+// aclk_query_t query = aclk_query_new(NODE_STATE_UPDATE);
39
+// node_instance_connection_t node_state_update = {
40
+// .hops = 1,
41
+// .live = 0,
42
+// .queryable = 1,
43
+// .session_id = aclk_session_newarch,
44
+// .node_id = node_id,
45
+// .capabilities = NULL};
46
+//
47
+// node_state_update.live = rrdhost_is_local(host) ? 1 : 0;
48
+// node_state_update.hops = rrdhost_ingestion_hops(host);
49
+// node_state_update.capabilities = aclk_get_node_instance_capas(host);
50
+// schedule_node_state_update(host, 5000);
51
+//
52
+// CLAIM_ID claim_id = claim_id_get();
53
+// node_state_update.claim_id = claim_id_is_set(claim_id) ? claim_id.str : NULL;
54
+// query->data.bin_payload.payload = generate_node_instance_connection(&query->data.bin_payload.size, &node_state_update);
55
+//
56
+// freez((void *)node_state_update.capabilities);
57
+//
58
+// query->data.bin_payload.msg_name = "UpdateNodeInstanceConnection";
59
+// query->data.bin_payload.topic = ACLK_TOPICID_NODE_CONN;
60
+//
61
+// aclk_add_job(query);
62
}
63
64
struct aclk_sync_config_s {
@@ -395,7 +395,7 @@ static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t query
395
break;
396
case SEND_NODE_INSTANCES:
397
worker_is_busy(UV_EVENT_SEND_NODE_INSTANCES);
398
- aclk_send_node_instances();
398
+ aclk_send_node_instances(config->client);
399
ok_to_send = false;
400
break;
401
case ALERT_START_STREAMING:
@@ -410,7 +410,7 @@ static void aclk_run_query(struct aclk_sync_config_s *config, aclk_query_t query
410
break;
411
case CREATE_NODE_INSTANCE:
412
worker_is_busy(UV_EVENT_CREATE_NODE_INSTANCE);
413
- create_node_instance_result_job(query->machine_guid, query->data.node_id);
413
+ create_node_instance_result_job(config->client, query->machine_guid, query->data.node_id);
414
ok_to_send = false;
415
break;
416
src/database/sqlite/sqlite_metadata.c
+7
-7
@@ -420,11 +420,10 @@ done:
420
421
#define SQL_UPDATE_NODE_ID "UPDATE node_instance SET node_id = @node_id WHERE host_id = @host_id"
422
423
-int sql_update_node_id(nd_uuid_t *host_id, nd_uuid_t *node_id)
423
+void sql_update_node_id(nd_uuid_t *host_id, nd_uuid_t *node_id)
424
{
425
sqlite3_stmt *res = NULL;
426
RRDHOST *host = NULL;
427
- int rc = 2;
427
428
char host_guid[GUID_LEN + 1];
429
uuid_unparse_lower(*host_id, host_guid);
@@ -435,25 +434,23 @@ int sql_update_node_id(nd_uuid_t *host_id, nd_uuid_t *node_id)
434
rrd_wrunlock();
435
436
if (!REQUIRE_DB(db_meta))
438
- return 1;
437
+ return;
438
439
if (!PREPARE_STATEMENT(db_meta, SQL_UPDATE_NODE_ID, &res))
441
- return 1;
440
+ return;
441
442
int param = 0;
443
SQLITE_BIND_FAIL(done, sqlite3_bind_blob(res, ++param, node_id, sizeof(*node_id), SQLITE_STATIC));
444
SQLITE_BIND_FAIL(done, sqlite3_bind_blob(res, ++param, host_id, sizeof(*host_id), SQLITE_STATIC));
445
446
param = 0;
448
- rc = execute_insert(res);
447
+ int rc = sqlite3_step_monitored(res);
448
if (unlikely(rc != SQLITE_DONE))
449
error_report("Failed to store node instance information, rc = %d", rc);
451
- rc = sqlite3_changes(db_meta);
450
451
done:
452
REPORT_BIND_FAIL(res, param);
453
SQLITE_FINALIZE(res);
456
- return rc - 1;
454
}
455
456
#define SQL_SELECT_NODE_ID "SELECT node_id FROM node_instance WHERE host_id = @host_id AND node_id IS NOT NULL"
@@ -534,6 +531,9 @@ struct node_instance_list *get_node_list(void)
531
while (sqlite3_step_monitored(res) == SQLITE_ROW)
532
row++;
533
534
+ if (row == 0)
535
+ return NULL;
536
+
537
if (sqlite3_reset(res) != SQLITE_OK) {
538
error_report("Failed to reset the prepared statement while fetching node instance information");
539
goto failed;
src/database/sqlite/sqlite_metadata.h
+1
-1
@@ -48,7 +48,7 @@ void vacuum_database(sqlite3 *database, const char *db_alias, int threshold, int
48
int sql_metadata_cache_stats(int op);
49
50
int get_node_id(nd_uuid_t *host_id, nd_uuid_t *node_id);
51
-int sql_update_node_id(nd_uuid_t *host_id, nd_uuid_t *node_id);
51
+void sql_update_node_id(nd_uuid_t *host_id, nd_uuid_t *node_id);
52
struct node_instance_list *get_node_list(void);
53
void sql_load_node_id(RRDHOST *host);
54