@cryptotaxi247 / netdata-1 / commits / a2852377d

Store and submit dimension delete messages for new cloud architecture (#11765)

* Enhance the dimension delete table and adjust the trigger to include chart_id and host_id * Add the aclk_process_dimension_deletion function * Change variable chart_name in aclk_upd_dimension_event (it is st->id from st.type dot st.id) * Process dimension deletion when retention updates are sent * Do not send charts if we don't have dimensions * Add check for uuid_parse return code

Stelios Fragkakis committed Nov 9, 2021 at 21:25 UTC a2852377d09811a18772f63de116c5b8068f2b59
4 files changed +82 -5
database/rrdset.c
+1 -1
@@ -1392,7 +1392,7 @@ void rrdset_done(RRDSET *st) {
1392
1393 #ifdef ENABLE_ACLK
1394 if (unlikely(!rrdset_flag_check(st, RRDSET_FLAG_ACLK))) {
1395 - if (st->counter_done >= RRDSET_MINIMUM_LIVE_COUNT) {
1395 + if (st->counter_done >= RRDSET_MINIMUM_LIVE_COUNT && st->dimensions) {
1396 if (likely(!queue_chart_to_aclk(st)))
1397 rrdset_flag_set(st, RRDSET_FLAG_ACLK);
1398 }
database/sqlite/sqlite_aclk.c
+11
@@ -11,6 +11,15 @@
11 #endif
12
13 const char *aclk_sync_config[] = {
14 + "CREATE TABLE IF NOT EXISTS dimension_delete (dimension_id blob, dimension_name text, chart_type_id text, "
15 + "dim_id blob, chart_id blob, host_id blob, date_created);",
16 +
17 + "CREATE INDEX IF NOT EXISTS ind_h1 ON dimension_delete (host_id);",
18 +
19 + "CREATE TRIGGER IF NOT EXISTS tr_dim_del AFTER DELETE ON dimension BEGIN INSERT INTO dimension_delete "
20 + "(dimension_id, dimension_name, chart_type_id, dim_id, chart_id, host_id, date_created)"
21 + " select old.id, old.name, c.type||\".\"||c.id, old.dim_id, old.chart_id, c.host_id, strftime('%s') FROM"
22 + " chart c WHERE c.chart_id = old.chart_id; END;",
23 NULL,
24 };
25
@@ -446,10 +455,12 @@ void aclk_database_worker(void *arg)
455 #ifdef ENABLE_NEW_CLOUD_PROTOCOL
456 case ACLK_DATABASE_DIM_DELETION:
457 debug(D_ACLK_SYNC,"Sending dimension deletion information %s", wc->uuid_str);
458 + aclk_process_dimension_deletion(wc, cmd);
459 break;
460 case ACLK_DATABASE_UPD_RETENTION:
461 debug(D_ACLK_SYNC,"Sending retention info for %s", wc->uuid_str);
462 aclk_update_retention(wc, cmd);
463 + aclk_process_dimension_deletion(wc, cmd);
464 break;
465 #endif
466
database/sqlite/sqlite_aclk_chart.c
+69 -4
@@ -170,24 +170,28 @@ int aclk_add_chart_event(struct aclk_database_worker_config *wc, struct aclk_dat
170 }
171
172 static inline int aclk_upd_dimension_event(struct aclk_database_worker_config *wc, char *claim_id, uuid_t *dim_uuid,
173 - const char *dim_id, const char *dim_name, const char *chart_name, time_t first_time, time_t last_time)
173 + const char *dim_id, const char *dim_name, const char *chart_type_id, time_t first_time, time_t last_time)
174 {
175 int rc = 0;
176 size_t size;
177
178 - if (unlikely(!dim_uuid || !dim_id || !dim_name || !chart_name))
178 + if (unlikely(!dim_uuid || !dim_id || !dim_name || !chart_type_id))
179 return 0;
180
181 struct chart_dimension_updated dim_payload;
182 memset(&dim_payload, 0, sizeof(dim_payload));
183
184 +#ifdef NETDATA_INTERNAL_CHECKS
185 if (!first_time)
185 - info("DEBUG: Deleting dimension [%s] [%s] [%s] [%s] [%s]", wc->node_id, claim_id, dim_id, dim_name, chart_name);
186 + info("Host %s (node %s) deleting dimension id=[%s] name=[%s] chart=[%s]",
187 + wc->host_guid, wc->node_id, dim_id, dim_name, chart_type_id);
188 +#endif
189 +
190 dim_payload.node_id = wc->node_id;
191 dim_payload.claim_id = claim_id;
192 dim_payload.name = dim_name;
193 dim_payload.id = dim_id;
190 - dim_payload.chart_id = chart_name;
194 + dim_payload.chart_id = chart_type_id;
195 dim_payload.created_at.tv_sec = first_time;
196 dim_payload.last_timestamp.tv_sec = last_time;
197 char *payload = generate_chart_dimension_updated(&size, &dim_payload);
@@ -197,6 +201,67 @@ static inline int aclk_upd_dimension_event(struct aclk_database_worker_config *w
201 return rc;
202 }
203
204 +void aclk_process_dimension_deletion(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd)
205 +{
206 + int rc = 0;
207 + sqlite3_stmt *res = NULL;
208 +
209 + if (!aclk_use_new_cloud_arch || !aclk_connected)
210 + return;
211 +
212 + if (unlikely(!db_meta))
213 + return;
214 +
215 + uuid_t host_id;
216 + if (uuid_parse(wc->host_guid, host_id))
217 + return;
218 +
219 + char *claim_id = is_agent_claimed();
220 + if (!claim_id)
221 + return;
222 +
223 + rc = sqlite3_prepare_v2(db_meta, "DELETE FROM dimension_delete where host_id = @host_id " \
224 + "RETURNING dimension_id, dimension_name, chart_type_id, dim_id LIMIT 10;", -1, &res, 0);
225 +
226 + if (rc != SQLITE_OK) {
227 + error_report("Failed to prepare statement when trying to delete dimension deletes");
228 + freez(claim_id);
229 + return;
230 + }
231 +
232 + rc = sqlite3_bind_blob(res, 1, &host_id , sizeof(host_id), SQLITE_STATIC);
233 + if (unlikely(rc != SQLITE_OK))
234 + goto bind_fail;
235 +
236 + unsigned count = 0;
237 + while (sqlite3_step(res) == SQLITE_ROW) {
238 + (void) aclk_upd_dimension_event(
239 + wc,
240 + claim_id,
241 + (uuid_t *)sqlite3_column_text(res, 3),
242 + (const char *)sqlite3_column_text(res, 0),
243 + (const char *)sqlite3_column_text(res, 1),
244 + (const char *)sqlite3_column_text(res, 2),
245 + 0,
246 + 0);
247 + count++;
248 + }
249 +
250 + if (count) {
251 + memset(&cmd, 0, sizeof(cmd));
252 + cmd.opcode = ACLK_DATABASE_DIM_DELETION;
253 + if (aclk_database_enq_cmd_noblock(wc, &cmd))
254 + info("Failed to queue a dimension deletion message");
255 + }
256 +
257 +bind_fail:
258 + rc = sqlite3_finalize(res);
259 + if (unlikely(rc != SQLITE_OK))
260 + error_report("Failed to finalize statement when adding dimension deletion events, rc = %d", rc);
261 + freez(claim_id);
262 + return;
263 +}
264 +
265 int aclk_add_dimension_event(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd)
266 {
267 int rc = 0;
database/sqlite/sqlite_aclk_chart.h
+1
@@ -32,5 +32,6 @@ void sql_check_rotation_state(struct aclk_database_worker_config *wc, struct acl
32 void sql_get_last_chart_sequence(struct aclk_database_worker_config *wc);
33 void aclk_receive_chart_reset(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd);
34 void aclk_receive_chart_ack(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd);
35 +void aclk_process_dimension_deletion(struct aclk_database_worker_config *wc, struct aclk_database_cmd cmd);
36 uint32_t sql_get_pending_count(struct aclk_database_worker_config *wc);
37 #endif //NETDATA_SQLITE_ACLK_CHART_H