| 1 | // SPDX-License-Identifier: GPL-3.0-or-later |
| 2 | |
| 3 | #include "query-internal.h" |
| 4 | |
| 5 | // ---------------------------------------------------------------------------- |
| 6 | // fill RRDR for the whole chart |
| 7 | |
| 8 | #ifdef NETDATA_INTERNAL_CHECKS |
| 9 | static void rrd2rrdr_log_request_response_metadata(RRDR *r |
| 10 | , RRDR_OPTIONS options __maybe_unused |
| 11 | , RRDR_TIME_GROUPING group_method |
| 12 | , bool aligned |
| 13 | , size_t group |
| 14 | , time_t resampling_time |
| 15 | , size_t resampling_group |
| 16 | , time_t after_wanted |
| 17 | , time_t after_requested |
| 18 | , time_t before_wanted |
| 19 | , time_t before_requested |
| 20 | , size_t points_requested |
| 21 | , size_t points_wanted |
| 22 | //, size_t after_slot |
| 23 | //, size_t before_slot |
| 24 | , const char *msg |
| 25 | ) { |
| 26 | |
| 27 | QUERY_TARGET *qt = r->internal.qt; |
| 28 | time_t first_entry_s = qt->db.first_time_s; |
| 29 | time_t last_entry_s = qt->db.last_time_s; |
| 30 | |
| 31 | internal_error( |
| 32 | true, |
| 33 | "rrd2rrdr() on %s update every %ld with %s grouping %s (group: %zu, resampling_time: %ld, resampling_group: %zu), " |
| 34 | "after (got: %ld, want: %ld, req: %ld, db: %ld), " |
| 35 | "before (got: %ld, want: %ld, req: %ld, db: %ld), " |
| 36 | "duration (got: %ld, want: %ld, req: %ld, db: %ld), " |
| 37 | "points (got: %zu, want: %zu, req: %zu), " |
| 38 | "%s" |
| 39 | , qt->id |
| 40 | , qt->window.query_granularity |
| 41 | |
| 42 | // grouping |
| 43 | , (aligned) ? "aligned" : "unaligned" |
| 44 | , time_grouping_id2txt(group_method) |
| 45 | , group |
| 46 | , resampling_time |
| 47 | , resampling_group |
| 48 | |
| 49 | // after |
| 50 | , r->view.after |
| 51 | , after_wanted |
| 52 | , after_requested |
| 53 | , first_entry_s |
| 54 | |
| 55 | // before |
| 56 | , r->view.before |
| 57 | , before_wanted |
| 58 | , before_requested |
| 59 | , last_entry_s |
| 60 | |
| 61 | // duration |
| 62 | , (long)(r->view.before - r->view.after + qt->window.query_granularity) |
| 63 | , (long)(before_wanted - after_wanted + qt->window.query_granularity) |
| 64 | , (long)before_requested - after_requested |
| 65 | , (long)((last_entry_s - first_entry_s) + qt->window.query_granularity) |
| 66 | |
| 67 | // points |
| 68 | , r->rows |
| 69 | , points_wanted |
| 70 | , points_requested |
| 71 | |
| 72 | // message |
| 73 | , msg |
| 74 | ); |
| 75 | } |
| 76 | #endif // NETDATA_INTERNAL_CHECKS |
| 77 | |
| 78 | // ---------------------------------------------------------------------------- |
| 79 | // query entry point |
| 80 | |
| 81 | RRDR *rrd2rrdr_legacy( |
| 82 | ONEWAYALLOC *owa, |
| 83 | RRDSET *st, size_t points, time_t after, time_t before, |
| 84 | RRDR_TIME_GROUPING group_method, time_t resampling_time, RRDR_OPTIONS options, const char *dimensions, |
| 85 | const char *group_options, time_t timeout_ms, size_t tier, QUERY_SOURCE query_source, |
| 86 | STORAGE_PRIORITY priority) { |
| 87 | |
| 88 | QUERY_TARGET_REQUEST qtr = { |
| 89 | .version = 1, |
| 90 | .st = st, |
| 91 | .points = points, |
| 92 | .after = after, |
| 93 | .before = before, |
| 94 | .time_group_method = group_method, |
| 95 | .resampling_time = resampling_time, |
| 96 | .options = options, |
| 97 | .dimensions = dimensions, |
| 98 | .time_group_options = group_options, |
| 99 | .timeout_ms = timeout_ms, |
| 100 | .tier = tier, |
| 101 | .query_source = query_source, |
| 102 | .priority = priority, |
| 103 | }; |
| 104 | |
| 105 | QUERY_TARGET *qt = query_target_create(&qtr); |
| 106 | RRDR *r = rrd2rrdr(owa, qt); |
| 107 | if(!r) { |
| 108 | query_target_release(qt); |
| 109 | return NULL; |
| 110 | } |
| 111 | |
| 112 | r->internal.release_with_rrdr_qt = qt; |
| 113 | return r; |
| 114 | } |
| 115 | |
| 116 | RRDR *rrd2rrdr(ONEWAYALLOC *owa, QUERY_TARGET *qt) { |
| 117 | if(!qt || !owa) |
| 118 | return NULL; |
| 119 | |
| 120 | // qt.window members are the WANTED ones. |
| 121 | // qt.request members are the REQUESTED ones. |
| 122 | |
| 123 | RRDR *r_tmp = rrd2rrdr_group_by_initialize(owa, qt); |
| 124 | if(!r_tmp) |
| 125 | return NULL; |
| 126 | |
| 127 | // the RRDR we group-by at |
| 128 | RRDR *r = (r_tmp->group_by.r) ? r_tmp->group_by.r : r_tmp; |
| 129 | |
| 130 | // the final RRDR to return to callers |
| 131 | RRDR *last_r = r_tmp; |
| 132 | while(last_r->group_by.r) |
| 133 | last_r = last_r->group_by.r; |
| 134 | |
| 135 | if(qt->window.relative) |
| 136 | last_r->view.flags |= RRDR_RESULT_FLAG_RELATIVE; |
| 137 | else |
| 138 | last_r->view.flags |= RRDR_RESULT_FLAG_ABSOLUTE; |
| 139 | |
| 140 | // ------------------------------------------------------------------------- |
| 141 | // assign the processor functions |
| 142 | rrdr_set_grouping_function(r_tmp, qt->window.time_group_method); |
| 143 | |
| 144 | // allocate any memory required by the grouping method |
| 145 | r_tmp->time_grouping.create(r_tmp, qt->window.time_group_options); |
| 146 | |
| 147 | // ------------------------------------------------------------------------- |
| 148 | // do the work for each dimension |
| 149 | |
| 150 | time_t max_after = 0, min_before = 0; |
| 151 | size_t max_rows = 0; |
| 152 | |
| 153 | long dimensions_used = 0, dimensions_nonzero = 0; |
| 154 | size_t last_db_points_read = 0; |
| 155 | size_t last_result_points_generated = 0; |
| 156 | |
| 157 | // internal_fatal(released_ops, "QUERY: released_ops should be NULL when the query starts"); |
| 158 | |
| 159 | query_progress_set_finish_line(qt->request.transaction, qt->query.used); |
| 160 | |
| 161 | QUERY_ENGINE_OPS **ops = NULL; |
| 162 | if(qt->query.used) |
| 163 | ops = onewayalloc_callocz(owa, qt->query.used, sizeof(QUERY_ENGINE_OPS *)); |
| 164 | |
| 165 | size_t capacity = MAX(netdata_conf_cpus() / 2, 4); |
| 166 | size_t max_queries_to_prepare = (qt->query.used > (capacity - 1)) ? (capacity - 1) : qt->query.used; |
| 167 | size_t queries_prepared = 0; |
| 168 | while(queries_prepared < max_queries_to_prepare) { |
| 169 | // preload another query |
| 170 | ops[queries_prepared] = rrd2rrdr_query_ops_prep(r_tmp, queries_prepared); |
| 171 | queries_prepared++; |
| 172 | } |
| 173 | |
| 174 | QUERY_NODE *last_qn = NULL; |
| 175 | usec_t last_ut = now_monotonic_usec(); |
| 176 | usec_t last_qn_ut = last_ut; |
| 177 | |
| 178 | for(size_t d = 0; d < qt->query.used ; d++) { |
| 179 | QUERY_METRIC *qm = query_metric(qt, d); |
| 180 | QUERY_DIMENSION *qd = query_dimension(qt, qm->link.query_dimension_id); |
| 181 | QUERY_INSTANCE *qi = query_instance(qt, qm->link.query_instance_id); |
| 182 | QUERY_CONTEXT *qc = query_context(qt, qm->link.query_context_id); |
| 183 | QUERY_NODE *qn = query_node(qt, qm->link.query_node_id); |
| 184 | |
| 185 | usec_t now_ut = last_ut; |
| 186 | if(qn != last_qn) { |
| 187 | if(last_qn) |
| 188 | last_qn->duration_ut = now_ut - last_qn_ut; |
| 189 | |
| 190 | last_qn = qn; |
| 191 | last_qn_ut = now_ut; |
| 192 | } |
| 193 | |
| 194 | if(queries_prepared < qt->query.used) { |
| 195 | // preload another query |
| 196 | ops[queries_prepared] = rrd2rrdr_query_ops_prep(r_tmp, queries_prepared); |
| 197 | queries_prepared++; |
| 198 | } |
| 199 | |
| 200 | size_t dim_in_rrdr_tmp = (r_tmp != r) ? 0 : d; |
| 201 | |
| 202 | // set the query target dimension options to rrdr |
| 203 | r_tmp->od[dim_in_rrdr_tmp] = qm->status; |
| 204 | |
| 205 | // reset the grouping for the new dimension |
| 206 | r_tmp->time_grouping.reset(r_tmp); |
| 207 | |
| 208 | if(ops[d]) { |
| 209 | rrd2rrdr_query_execute(r_tmp, dim_in_rrdr_tmp, ops[d]); |
| 210 | r_tmp->od[dim_in_rrdr_tmp] |= RRDR_DIMENSION_QUERIED; |
| 211 | |
| 212 | now_ut = now_monotonic_usec(); |
| 213 | qm->duration_ut = now_ut - last_ut; |
| 214 | last_ut = now_ut; |
| 215 | |
| 216 | if(r_tmp != r) { |
| 217 | // copy back whatever got updated from the temporary r |
| 218 | |
| 219 | // the query updates RRDR_DIMENSION_NONZERO |
| 220 | qm->status = r_tmp->od[dim_in_rrdr_tmp]; |
| 221 | |
| 222 | // the query updates these |
| 223 | r->view.min = r_tmp->view.min; |
| 224 | r->view.max = r_tmp->view.max; |
| 225 | r->view.after = r_tmp->view.after; |
| 226 | r->view.before = r_tmp->view.before; |
| 227 | r->rows = r_tmp->rows; |
| 228 | |
| 229 | rrd2rrdr_group_by_add_metric(r, qm->grouped_as.first_slot, r_tmp, dim_in_rrdr_tmp, |
| 230 | qt->request.group_by[0].aggregation, &qm->query_points, 0); |
| 231 | } |
| 232 | |
| 233 | rrd2rrdr_query_ops_release(ops[d]); // reuse this ops allocation |
| 234 | ops[d] = NULL; |
| 235 | |
| 236 | qi->metrics.queried++; |
| 237 | qc->metrics.queried++; |
| 238 | qn->metrics.queried++; |
| 239 | |
| 240 | qd->status |= QUERY_STATUS_QUERIED; |
| 241 | qm->status |= RRDR_DIMENSION_QUERIED; |
| 242 | |
| 243 | if(qt->request.version >= 2) { |
| 244 | // we need to make the query points positive now |
| 245 | // since we will aggregate it across multiple dimensions |
| 246 | storage_point_make_positive(qm->query_points); |
| 247 | storage_point_merge_to(qi->query_points, qm->query_points); |
| 248 | storage_point_merge_to(qc->query_points, qm->query_points); |
| 249 | storage_point_merge_to(qn->query_points, qm->query_points); |
| 250 | storage_point_merge_to(qt->query_points, qm->query_points); |
| 251 | } |
| 252 | } |
| 253 | else { |
| 254 | qi->metrics.failed++; |
| 255 | qc->metrics.failed++; |
| 256 | qn->metrics.failed++; |
| 257 | |
| 258 | qd->status |= QUERY_STATUS_FAILED; |
| 259 | qm->status |= RRDR_DIMENSION_FAILED; |
| 260 | |
| 261 | continue; |
| 262 | } |
| 263 | |
| 264 | pulse_queries_rrdr_query_completed( |
| 265 | 1, |
| 266 | r_tmp->stats.db_points_read - last_db_points_read, |
| 267 | r_tmp->stats.result_points_generated - last_result_points_generated, |
| 268 | qt->request.query_source); |
| 269 | |
| 270 | last_db_points_read = r_tmp->stats.db_points_read; |
| 271 | last_result_points_generated = r_tmp->stats.result_points_generated; |
| 272 | |
| 273 | if(qm->status & RRDR_DIMENSION_NONZERO) |
| 274 | dimensions_nonzero++; |
| 275 | |
| 276 | // verify all dimensions are aligned |
| 277 | if(unlikely(!dimensions_used)) { |
| 278 | min_before = r->view.before; |
| 279 | max_after = r->view.after; |
| 280 | max_rows = r->rows; |
| 281 | } |
| 282 | else { |
| 283 | if(r->view.after != max_after) { |
| 284 | internal_error(true, "QUERY: 'after' mismatch between dimensions for chart '%s': max is %zu, dimension '%s' has %zu", |
| 285 | rrdinstance_acquired_id(qi->ria), (size_t)max_after, rrdmetric_acquired_id(qd->rma), (size_t)r->view.after); |
| 286 | |
| 287 | r->view.after = (r->view.after > max_after) ? r->view.after : max_after; |
| 288 | } |
| 289 | |
| 290 | if(r->view.before != min_before) { |
| 291 | internal_error(true, "QUERY: 'before' mismatch between dimensions for chart '%s': max is %zu, dimension '%s' has %zu", |
| 292 | rrdinstance_acquired_id(qi->ria), (size_t)min_before, rrdmetric_acquired_id(qd->rma), (size_t)r->view.before); |
| 293 | |
| 294 | r->view.before = (r->view.before < min_before) ? r->view.before : min_before; |
| 295 | } |
| 296 | |
| 297 | if(r->rows != max_rows) { |
| 298 | internal_error(true, "QUERY: 'rows' mismatch between dimensions for chart '%s': max is %zu, dimension '%s' has %zu", |
| 299 | rrdinstance_acquired_id(qi->ria), (size_t)max_rows, rrdmetric_acquired_id(qd->rma), (size_t)r->rows); |
| 300 | |
| 301 | r->rows = (r->rows > max_rows) ? r->rows : max_rows; |
| 302 | } |
| 303 | } |
| 304 | |
| 305 | dimensions_used++; |
| 306 | |
| 307 | bool cancel = false; |
| 308 | if (qt->request.interrupt_callback && qt->request.interrupt_callback(qt->request.interrupt_callback_data)) { |
| 309 | cancel = true; |
| 310 | nd_log(NDLS_ACCESS, NDLP_NOTICE, "QUERY INTERRUPTED"); |
| 311 | } |
| 312 | |
| 313 | if (qt->request.timeout_ms && ((NETDATA_DOUBLE)(now_ut - qt->timings.received_ut) / 1000.0) > (NETDATA_DOUBLE)qt->request.timeout_ms) { |
| 314 | cancel = true; |
| 315 | nd_log(NDLS_ACCESS, NDLP_WARNING, "QUERY CANCELED RUNTIME EXCEEDED %0.2f ms (LIMIT %lld ms)", |
| 316 | (NETDATA_DOUBLE)(now_ut - qt->timings.received_ut) / 1000.0, (long long)qt->request.timeout_ms); |
| 317 | } |
| 318 | |
| 319 | if(cancel) { |
| 320 | r->view.flags |= RRDR_RESULT_FLAG_CANCEL; |
| 321 | |
| 322 | for(size_t i = d + 1; i < queries_prepared ; i++) { |
| 323 | if(ops[i]) { |
| 324 | query_planer_finalize_remaining_plans(ops[i]); |
| 325 | rrd2rrdr_query_ops_release(ops[i]); |
| 326 | ops[i] = NULL; |
| 327 | } |
| 328 | } |
| 329 | |
| 330 | break; |
| 331 | } |
| 332 | else |
| 333 | query_progress_done_step(qt->request.transaction, 1); |
| 334 | } |
| 335 | |
| 336 | // free all resources used by the grouping method |
| 337 | r_tmp->time_grouping.free(r_tmp); |
| 338 | |
| 339 | // get the final RRDR to send to the caller |
| 340 | r = rrd2rrdr_group_by_finalize(r_tmp); |
| 341 | |
| 342 | // apply cardinality limit if requested |
| 343 | r = rrd2rrdr_cardinality_limit(r); |
| 344 | |
| 345 | #ifdef NETDATA_INTERNAL_CHECKS |
| 346 | if (dimensions_used && !(r->view.flags & RRDR_RESULT_FLAG_CANCEL)) { |
| 347 | if(r->internal.log) |
| 348 | rrd2rrdr_log_request_response_metadata(r, qt->window.options, qt->window.time_group_method, qt->window.aligned, qt->window.group, qt->request.resampling_time, qt->window.resampling_group, |
| 349 | qt->window.after, qt->request.after, qt->window.before, qt->request.before, |
| 350 | qt->request.points, qt->window.points, /*after_slot, before_slot,*/ |
| 351 | r->internal.log); |
| 352 | |
| 353 | if(r->rows != qt->window.points) |
| 354 | rrd2rrdr_log_request_response_metadata(r, qt->window.options, qt->window.time_group_method, qt->window.aligned, qt->window.group, qt->request.resampling_time, qt->window.resampling_group, |
| 355 | qt->window.after, qt->request.after, qt->window.before, qt->request.before, |
| 356 | qt->request.points, qt->window.points, /*after_slot, before_slot,*/ |
| 357 | "got 'points' is not wanted 'points'"); |
| 358 | |
| 359 | if(qt->window.aligned && (r->view.before % query_view_update_every(qt)) != 0) |
| 360 | rrd2rrdr_log_request_response_metadata(r, qt->window.options, qt->window.time_group_method, qt->window.aligned, qt->window.group, qt->request.resampling_time, qt->window.resampling_group, |
| 361 | qt->window.after, qt->request.after, qt->window.before, qt->request.before, |
| 362 | qt->request.points, qt->window.points, /*after_slot, before_slot,*/ |
| 363 | "'before' is not aligned but alignment is required"); |
| 364 | |
| 365 | // 'after' should not be aligned, since we start inside the first group |
| 366 | //if(qt->window.aligned && (r->after % group) != 0) |
| 367 | // rrd2rrdr_log_request_response_metadata(r, qt->window.options, qt->window.group_method, qt->window.aligned, qt->window.group, qt->request.resampling_time, qt->window.resampling_group, qt->window.after, after_requested, before_wanted, before_requested, points_requested, points_wanted, after_slot, before_slot, "'after' is not aligned but alignment is required"); |
| 368 | |
| 369 | if(r->view.before != qt->window.before) |
| 370 | rrd2rrdr_log_request_response_metadata(r, qt->window.options, qt->window.time_group_method, qt->window.aligned, qt->window.group, qt->request.resampling_time, qt->window.resampling_group, |
| 371 | qt->window.after, qt->request.after, qt->window.before, qt->request.before, |
| 372 | qt->request.points, qt->window.points, /*after_slot, before_slot,*/ |
| 373 | "chart is not aligned to requested 'before'"); |
| 374 | |
| 375 | if(r->view.before != qt->window.before) |
| 376 | rrd2rrdr_log_request_response_metadata(r, qt->window.options, qt->window.time_group_method, qt->window.aligned, qt->window.group, qt->request.resampling_time, qt->window.resampling_group, |
| 377 | qt->window.after, qt->request.after, qt->window.before, qt->request.before, |
| 378 | qt->request.points, qt->window.points, /*after_slot, before_slot,*/ |
| 379 | "got 'before' is not wanted 'before'"); |
| 380 | |
| 381 | // reported 'after' varies, depending on group |
| 382 | if(r->view.after != qt->window.after) |
| 383 | rrd2rrdr_log_request_response_metadata(r, qt->window.options, qt->window.time_group_method, qt->window.aligned, qt->window.group, qt->request.resampling_time, qt->window.resampling_group, |
| 384 | qt->window.after, qt->request.after, qt->window.before, qt->request.before, |
| 385 | qt->request.points, qt->window.points, /*after_slot, before_slot,*/ |
| 386 | "got 'after' is not wanted 'after'"); |
| 387 | |
| 388 | } |
| 389 | #endif |
| 390 | |
| 391 | // free the query pipelining ops |
| 392 | for(size_t d = 0; d < qt->query.used ; d++) { |
| 393 | rrd2rrdr_query_ops_release(ops[d]); |
| 394 | ops[d] = NULL; |
| 395 | } |
| 396 | rrd2rrdr_query_ops_freeall(r); |
| 397 | // internal_fatal(released_ops, "QUERY: released_ops should be NULL when the query ends"); |
| 398 | |
| 399 | onewayalloc_freez(owa, ops); |
| 400 | |
| 401 | if(likely(dimensions_used && (qt->window.options & RRDR_OPTION_NONZERO) && !dimensions_nonzero)) |
| 402 | // when all the dimensions are zero, we should return all of them |
| 403 | qt->window.options &= ~RRDR_OPTION_NONZERO; |
| 404 | |
| 405 | qt->timings.executed_ut = now_monotonic_usec(); |
| 406 | |
| 407 | return r; |
| 408 | } |