master
c 408 lines 16.7 KB
Raw
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 }