master
c 993 lines 41.8 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 #include "rrdset-collection.h"
4 #include "rrddim-collection.h"
5
6 time_t rrdset_set_update_every_s(RRDSET *st, time_t update_every_s) {
7 if(unlikely(update_every_s == st->update_every))
8 return st->update_every;
9
10 internal_error(true, "RRDSET '%s' switching update every from %d to %d",
11 rrdset_id(st), (int)st->update_every, (int)update_every_s);
12
13 time_t prev_update_every_s = (time_t) st->update_every;
14 st->update_every = (int) update_every_s;
15
16 // switch update every to the storage engine
17 RRDDIM *rd;
18 rrddim_foreach_read(rd, st) {
19 for (size_t tier = 0; tier < nd_profile.storage_tiers; tier++) {
20 if (rd->tiers[tier].sch)
21 storage_engine_store_change_collection_frequency(
22 rd->tiers[tier].sch,
23 (int)(st->rrdhost->db[tier].tier_grouping * st->update_every));
24 }
25 }
26 rrddim_foreach_done(rd);
27
28 return prev_update_every_s;
29 }
30
31 void rrdset_finalize_collection(RRDSET *st, bool dimensions_too) {
32 ND_LOG_STACK lgs[] = {
33 ND_LOG_FIELD_TXT(NDF_NIDL_NODE, rrdhost_hostname(st->rrdhost)),
34 ND_LOG_FIELD_TXT(NDF_NIDL_CONTEXT, rrdset_context(st)),
35 ND_LOG_FIELD_TXT(NDF_NIDL_INSTANCE, rrdset_name(st)),
36 ND_LOG_FIELD_END(),
37 };
38 ND_LOG_STACK_PUSH(lgs);
39
40 RRDHOST *host = st->rrdhost;
41
42 rrdset_flag_set(st, RRDSET_FLAG_COLLECTION_FINISHED);
43
44 if(dimensions_too) {
45 RRDDIM *rd;
46 rrddim_foreach_read(rd, st)
47 rrddim_finalize_collection_and_check_retention(rd);
48 rrddim_foreach_done(rd);
49 }
50
51 for(size_t tier = 0; tier < nd_profile.storage_tiers; tier++) {
52 STORAGE_ENGINE *eng = st->rrdhost->db[tier].eng;
53 if(!eng) continue;
54
55 if(st->smg[tier]) {
56 storage_engine_metrics_group_release(eng->seb, host->db[tier].si, st->smg[tier]);
57 st->smg[tier] = NULL;
58 }
59 }
60
61 rrdset_pluginsd_receive_unslot_and_cleanup(st);
62 }
63
64 // ----------------------------------------------------------------------------
65 // RRDSET - reset a chart
66
67 static void rrdset_collection_reset(RRDSET *st) {
68 netdata_log_debug(D_RRD_CALLS, "rrdset_collection_reset() %s", rrdset_name(st));
69
70 st->last_collected_time.tv_sec = 0;
71 st->last_collected_time.tv_usec = 0;
72 st->last_updated.tv_sec = 0;
73 st->last_updated.tv_usec = 0;
74 st->db.current_entry = 0;
75 st->counter = 0;
76 st->counter_done = 0;
77
78 RRDDIM *rd;
79 rrddim_foreach_read(rd, st) {
80 rd->collector.last_collected_time.tv_sec = 0;
81 rd->collector.last_collected_time.tv_usec = 0;
82 rd->collector.counter = 0;
83
84 for(size_t tier = 0; tier < nd_profile.storage_tiers;tier++)
85 storage_engine_store_flush(rd->tiers[tier].sch);
86 }
87 rrddim_foreach_done(rd);
88 }
89
90 // ----------------------------------------------------------------------------
91 // RRDSET - data collection iteration control
92
93 static inline void last_collected_time_align(RRDSET *st) {
94 st->last_collected_time.tv_sec -= st->last_collected_time.tv_sec % st->update_every;
95
96 if(unlikely(rrdset_flag_check(st, RRDSET_FLAG_STORE_FIRST)))
97 st->last_collected_time.tv_usec = 0;
98 else
99 st->last_collected_time.tv_usec = 500000;
100 }
101
102 static inline void last_updated_time_align(RRDSET *st) {
103 st->last_updated.tv_sec -= st->last_updated.tv_sec % st->update_every;
104 st->last_updated.tv_usec = 0;
105 }
106
107 void rrdset_timed_next(RRDSET *st, struct timeval now, usec_t duration_since_last_update) {
108 #ifdef NETDATA_INTERNAL_CHECKS
109 char *discard_reason = NULL;
110 usec_t discarded = duration_since_last_update;
111 #endif
112
113 if(unlikely(rrdset_flag_check(st, RRDSET_FLAG_SYNC_CLOCK))) {
114 // the chart needs to be re-synced to current time
115 rrdset_flag_clear(st, RRDSET_FLAG_SYNC_CLOCK);
116
117 // discard the duration supplied
118 duration_since_last_update = 0;
119
120 #ifdef NETDATA_INTERNAL_CHECKS
121 if(!discard_reason) discard_reason = "SYNC CLOCK FLAG";
122 #endif
123 }
124
125 if(unlikely(!st->last_collected_time.tv_sec)) {
126 // the first entry
127 duration_since_last_update = st->update_every * USEC_PER_SEC;
128 #ifdef NETDATA_INTERNAL_CHECKS
129 if(!discard_reason) discard_reason = "FIRST DATA COLLECTION";
130 #endif
131 }
132 else if(unlikely(!duration_since_last_update)) {
133 // no dt given by the plugin
134 duration_since_last_update = dt_usec(&now, &st->last_collected_time);
135 #ifdef NETDATA_INTERNAL_CHECKS
136 if(!discard_reason) discard_reason = "NO USEC GIVEN BY COLLECTOR";
137 #endif
138 }
139 else {
140 // microseconds has the time since the last collection
141 susec_t since_last_usec = dt_usec_signed(&now, &st->last_collected_time);
142
143 if(unlikely(since_last_usec < 0)) {
144 // oops! the database is in the future
145 #ifdef NETDATA_INTERNAL_CHECKS
146 netdata_log_info("RRD database for chart '%s' on host '%s' is %0.5" NETDATA_DOUBLE_MODIFIER
147 " secs in the future (counter #%u, update #%u). Adjusting it to current time."
148 , rrdset_id(st)
149 , rrdhost_hostname(st->rrdhost)
150 , (NETDATA_DOUBLE)-since_last_usec / USEC_PER_SEC
151 , st->counter
152 , st->counter_done
153 );
154 #endif
155
156 duration_since_last_update = 0;
157
158 #ifdef NETDATA_INTERNAL_CHECKS
159 if(!discard_reason) discard_reason = "COLLECTION TIME IN FUTURE";
160 #endif
161 }
162 else if(unlikely((usec_t)since_last_usec > (usec_t)(st->update_every * 5 * USEC_PER_SEC))) {
163 // oops! the database is too far behind
164 #ifdef NETDATA_INTERNAL_CHECKS
165 netdata_log_info("RRD database for chart '%s' on host '%s' is %0.5" NETDATA_DOUBLE_MODIFIER
166 " secs in the past (counter #%u, update #%u). Adjusting it to current time.",
167 rrdset_id(st), rrdhost_hostname(st->rrdhost), (NETDATA_DOUBLE)since_last_usec / USEC_PER_SEC,
168 st->counter, st->counter_done);
169 #endif
170
171 duration_since_last_update = (usec_t)since_last_usec;
172
173 #ifdef NETDATA_INTERNAL_CHECKS
174 if(!discard_reason) discard_reason = "COLLECTION TIME TOO FAR IN THE PAST";
175 #endif
176 }
177
178 #ifdef NETDATA_INTERNAL_CHECKS
179 if(since_last_usec > 0 && (susec_t) duration_since_last_update < since_last_usec) {
180 static __thread susec_t min_delta = USEC_PER_SEC * 3600, permanent_min_delta = 0;
181 static __thread time_t last_time_s = 0;
182
183 // the first time initialize it so that it will make the check later
184 if(last_time_s == 0) last_time_s = now.tv_sec + 60;
185
186 susec_t delta = since_last_usec - (susec_t) duration_since_last_update;
187 if(delta < min_delta) min_delta = delta;
188
189 if(now.tv_sec >= last_time_s + 60) {
190 last_time_s = now.tv_sec;
191
192 if(min_delta > permanent_min_delta) {
193 netdata_log_info("MINIMUM MICROSECONDS DELTA of thread %d increased from %"PRIi64" to %"PRIi64" (+%"PRIi64")", gettid_cached(), permanent_min_delta, min_delta, min_delta - permanent_min_delta);
194 permanent_min_delta = min_delta;
195 }
196
197 min_delta = USEC_PER_SEC * 3600;
198 }
199 }
200 #endif
201 }
202
203 netdata_log_debug(D_RRD_CALLS, "rrdset_timed_next() for chart %s with duration since last update %"PRIu64" usec", rrdset_name(st), duration_since_last_update);
204 rrdset_debug(st, "NEXT: %"PRIu64" microseconds", duration_since_last_update);
205
206 internal_error(discarded && discarded != duration_since_last_update,
207 "host '%s', chart '%s': discarded data collection time of %"PRIu64" usec, "
208 "replaced with %"PRIu64" usec, reason: '%s'"
209 , rrdhost_hostname(st->rrdhost)
210 , rrdset_id(st)
211 , discarded
212 , duration_since_last_update
213 , discard_reason?discard_reason:"UNDEFINED"
214 );
215
216 st->usec_since_last_update = duration_since_last_update;
217 }
218
219 inline void rrdset_next_usec_unfiltered(RRDSET *st, usec_t duration_since_last_update) {
220 if(unlikely(!st->last_collected_time.tv_sec || !duration_since_last_update || (rrdset_flag_check(st, RRDSET_FLAG_SYNC_CLOCK)))) {
221 // call the full next_usec() function
222 rrdset_next_usec(st, duration_since_last_update);
223 return;
224 }
225
226 st->usec_since_last_update = duration_since_last_update;
227 }
228
229 inline void rrdset_next_usec(RRDSET *st, usec_t duration_since_last_update) {
230 struct timeval now;
231
232 now_realtime_timeval(&now);
233 rrdset_timed_next(st, now, duration_since_last_update);
234 }
235
236 // ----------------------------------------------------------------------------
237 // RRDSET - process the collected values for all dimensions of a chart
238
239 static inline usec_t rrdset_init_last_collected_time(RRDSET *st, struct timeval now) {
240 st->last_collected_time = now;
241 last_collected_time_align(st);
242
243 usec_t last_collect_ut = st->last_collected_time.tv_sec * USEC_PER_SEC + st->last_collected_time.tv_usec;
244
245 rrdset_debug(st, "initialized last collected time to %0.3" NETDATA_DOUBLE_MODIFIER, (NETDATA_DOUBLE)last_collect_ut / USEC_PER_SEC);
246
247 return last_collect_ut;
248 }
249
250 static inline usec_t rrdset_update_last_collected_time(RRDSET *st) {
251 usec_t last_collect_ut = st->last_collected_time.tv_sec * USEC_PER_SEC + st->last_collected_time.tv_usec;
252 usec_t ut = last_collect_ut + st->usec_since_last_update;
253 st->last_collected_time.tv_sec = (time_t) (ut / USEC_PER_SEC);
254 st->last_collected_time.tv_usec = (suseconds_t) (ut % USEC_PER_SEC);
255
256 rrdset_debug(st, "updated last collected time to %0.3" NETDATA_DOUBLE_MODIFIER, (NETDATA_DOUBLE)last_collect_ut / USEC_PER_SEC);
257
258 return last_collect_ut;
259 }
260
261 static inline void rrdset_init_last_updated_time(RRDSET *st) {
262 // copy the last collected time to last updated time
263 st->last_updated.tv_sec = st->last_collected_time.tv_sec;
264 st->last_updated.tv_usec = st->last_collected_time.tv_usec;
265
266 if(rrdset_flag_check(st, RRDSET_FLAG_STORE_FIRST))
267 st->last_updated.tv_sec -= st->update_every;
268
269 last_updated_time_align(st);
270 }
271
272 __thread size_t rrdset_done_statistics_points_stored_per_tier[RRD_STORAGE_TIERS];
273
274 // caching of dimensions rrdset_done() and rrdset_done_interpolate() loop through
275 struct rda_item {
276 const DICTIONARY_ITEM *item;
277 RRDDIM *rd;
278 bool reset_or_overflow;
279 };
280
281 static __thread struct rda_item *thread_rda = NULL;
282 static __thread size_t thread_rda_entries = 0;
283
284 static struct rda_item *rrdset_thread_rda_get(size_t *dimensions) {
285
286 if(unlikely(!thread_rda || (*dimensions) > thread_rda_entries)) {
287 size_t old_mem = thread_rda_entries * sizeof(struct rda_item);
288 freez(thread_rda);
289 thread_rda_entries = *dimensions;
290 size_t new_mem = thread_rda_entries * sizeof(struct rda_item);
291 thread_rda = mallocz(new_mem);
292
293 __atomic_add_fetch(&netdata_buffers_statistics.rrdset_done_rda_size, new_mem - old_mem, __ATOMIC_RELAXED);
294 }
295
296 *dimensions = thread_rda_entries;
297 return thread_rda;
298 }
299
300 void rrdset_thread_rda_free(void) {
301 __atomic_sub_fetch(&netdata_buffers_statistics.rrdset_done_rda_size, thread_rda_entries * sizeof(struct rda_item), __ATOMIC_RELAXED);
302
303 freez(thread_rda);
304 thread_rda = NULL;
305 thread_rda_entries = 0;
306 }
307
308 static inline size_t rrdset_done_interpolate(
309 RRDSET_STREAM_BUFFER *rsb
310 , RRDSET *st
311 , struct rda_item *rda_base
312 , size_t rda_slots
313 , usec_t update_every_ut
314 , usec_t last_stored_ut
315 , usec_t next_store_ut
316 , usec_t last_collect_ut
317 , usec_t now_collect_ut
318 , char store_this_entry
319 ) {
320 RRDDIM *rd;
321
322 size_t stored_entries = 0; // the number of entries we have stored in the db, during this call to rrdset_done()
323
324 usec_t first_ut = last_stored_ut, last_ut = 0;
325 (void)first_ut;
326
327 ssize_t iterations = (ssize_t)((now_collect_ut - last_stored_ut) / (update_every_ut));
328 if((now_collect_ut % (update_every_ut)) == 0) iterations++;
329
330 size_t counter = st->counter;
331 long current_entry = st->db.current_entry;
332
333 for( ; next_store_ut <= now_collect_ut ; last_collect_ut = next_store_ut, next_store_ut += update_every_ut, iterations-- ) {
334
335 internal_error(iterations < 0,
336 "RRDSET: '%s': iterations calculation wrapped! "
337 "first_ut = %"PRIu64", last_stored_ut = %"PRIu64", next_store_ut = %"PRIu64", now_collect_ut = %"PRIu64""
338 , rrdset_id(st)
339 , first_ut
340 , last_stored_ut
341 , next_store_ut
342 , now_collect_ut
343 );
344
345 rrdset_debug(st, "last_stored_ut = %0.3" NETDATA_DOUBLE_MODIFIER " (last updated time)", (NETDATA_DOUBLE)last_stored_ut/USEC_PER_SEC);
346 rrdset_debug(st, "next_store_ut = %0.3" NETDATA_DOUBLE_MODIFIER " (next interpolation point)", (NETDATA_DOUBLE)next_store_ut/USEC_PER_SEC);
347
348 last_ut = next_store_ut;
349
350 ml_chart_update_begin(st);
351
352 struct rda_item *rda;
353 size_t dim_id;
354 for(dim_id = 0, rda = rda_base ; dim_id < rda_slots ; ++dim_id, ++rda) {
355 rd = rda->rd;
356 if(unlikely(!rd)) continue;
357
358 SN_FLAGS storage_flags = SN_DEFAULT_FLAGS;
359
360 if (rda->reset_or_overflow)
361 storage_flags |= SN_FLAG_RESET;
362
363 NETDATA_DOUBLE new_value;
364
365 switch(rd->algorithm) {
366 case RRD_ALGORITHM_INCREMENTAL:
367 new_value = (NETDATA_DOUBLE)
368 ( rd->collector.calculated_value
369 * (NETDATA_DOUBLE)(next_store_ut - last_collect_ut)
370 / (NETDATA_DOUBLE)(now_collect_ut - last_collect_ut)
371 );
372
373 rrdset_debug(st, "%s: CALC2 INC " NETDATA_DOUBLE_FORMAT " = "
374 NETDATA_DOUBLE_FORMAT
375 " * (%"PRIu64" - %"PRIu64")"
376 " / (%"PRIu64" - %"PRIu64""
377 , rrddim_name(rd)
378 , new_value
379 , rd->collector.calculated_value
380 , next_store_ut, last_collect_ut
381 , now_collect_ut, last_collect_ut
382 );
383
384 rd->collector.calculated_value -= new_value;
385 new_value += rd->collector.last_calculated_value;
386 rd->collector.last_calculated_value = 0;
387 new_value /= (NETDATA_DOUBLE)st->update_every;
388
389 if(unlikely(next_store_ut - last_stored_ut < update_every_ut)) {
390
391 rrdset_debug(st, "%s: COLLECTION POINT IS SHORT " NETDATA_DOUBLE_FORMAT " - EXTRAPOLATING",
392 rrddim_name(rd)
393 , (NETDATA_DOUBLE)(next_store_ut - last_stored_ut)
394 );
395
396 new_value = new_value * (NETDATA_DOUBLE)(st->update_every * USEC_PER_SEC) / (NETDATA_DOUBLE)(next_store_ut - last_stored_ut);
397 }
398 break;
399
400 case RRD_ALGORITHM_ABSOLUTE:
401 case RRD_ALGORITHM_PCENT_OVER_ROW_TOTAL:
402 case RRD_ALGORITHM_PCENT_OVER_DIFF_TOTAL:
403 default:
404 if(iterations == 1) {
405 // this is the last iteration
406 // do not interpolate
407 // just show the calculated value
408
409 new_value = rd->collector.calculated_value;
410 }
411 else {
412 // we have missed an update
413 // interpolate in the middle values
414
415 new_value = (NETDATA_DOUBLE)
416 ( ( (rd->collector.calculated_value - rd->collector.last_calculated_value)
417 * (NETDATA_DOUBLE)(next_store_ut - last_collect_ut)
418 / (NETDATA_DOUBLE)(now_collect_ut - last_collect_ut)
419 )
420 + rd->collector.last_calculated_value
421 );
422
423 rrdset_debug(st, "%s: CALC2 DEF " NETDATA_DOUBLE_FORMAT " = ((("
424 "(" NETDATA_DOUBLE_FORMAT " - " NETDATA_DOUBLE_FORMAT ")"
425 " * %"PRIu64""
426 " / %"PRIu64") + " NETDATA_DOUBLE_FORMAT, rrddim_name(rd)
427 , new_value
428 , rd->collector.calculated_value, rd->collector.last_calculated_value
429 , (next_store_ut - first_ut)
430 , (now_collect_ut - first_ut), rd->collector.last_calculated_value
431 );
432 }
433 break;
434 }
435
436 time_t current_time_s = (time_t) (next_store_ut / USEC_PER_SEC);
437
438 if(unlikely(!store_this_entry)) {
439 (void) ml_dimension_is_anomalous(rd, current_time_s, 0, false);
440
441 if(rsb->wb && rsb->v2)
442 stream_send_rrddim_metrics_v2(rsb, rd, next_store_ut, NAN, SN_FLAG_NONE);
443
444 rrddim_store_metric(rd, next_store_ut, NAN, SN_FLAG_NONE);
445 continue;
446 }
447
448 if(likely(rrddim_check_updated(rd) && rd->collector.counter > 1 && iterations < gap_when_lost_iterations_above)) {
449 uint32_t dim_storage_flags = storage_flags;
450
451 if (ml_dimension_is_anomalous(rd, current_time_s, new_value, true)) {
452 // clear anomaly bit: 0 -> is anomalous, 1 -> not anomalous
453 dim_storage_flags &= ~((storage_number)SN_FLAG_NOT_ANOMALOUS);
454 }
455
456 if(rsb->wb && rsb->v2)
457 stream_send_rrddim_metrics_v2(rsb, rd, next_store_ut, new_value, dim_storage_flags);
458
459 rrddim_store_metric(rd, next_store_ut, new_value, dim_storage_flags);
460 rd->collector.last_stored_value = new_value;
461 }
462 else {
463 (void) ml_dimension_is_anomalous(rd, current_time_s, 0, false);
464
465 rrdset_debug(st, "%s: STORE[%ld] = NON EXISTING ", rrddim_name(rd), current_entry);
466
467 if(rsb->wb && rsb->v2)
468 stream_send_rrddim_metrics_v2(rsb, rd, next_store_ut, NAN, SN_FLAG_NONE);
469
470 rrddim_store_metric(rd, next_store_ut, NAN, SN_FLAG_NONE);
471 rd->collector.last_stored_value = NAN;
472 }
473
474 stored_entries++;
475 }
476
477 ml_chart_update_end(st);
478
479 st->counter = ++counter;
480 st->db.current_entry = current_entry = ((current_entry + 1) >= st->db.entries) ? 0 : current_entry + 1;
481
482 st->last_updated.tv_sec = (time_t) (last_ut / USEC_PER_SEC);
483 st->last_updated.tv_usec = 0;
484
485 last_stored_ut = next_store_ut;
486 }
487
488 /*
489 st->counter = counter;
490 st->current_entry = current_entry;
491
492 if(likely(last_ut)) {
493 st->last_updated.tv_sec = (time_t) (last_ut / USEC_PER_SEC);
494 st->last_updated.tv_usec = 0;
495 }
496 */
497
498 return stored_entries;
499 }
500
501 void rrdset_done(RRDSET *st) {
502 struct timeval now;
503
504 now_realtime_timeval(&now);
505 rrdset_timed_done(st, now, /* pending_rrdset_next = */ st->counter_done != 0);
506 }
507
508 void rrdset_timed_done(RRDSET *st, struct timeval now, bool pending_rrdset_next) {
509 if(unlikely(!service_running(SERVICE_COLLECTORS))) return;
510
511 RRDSET_STREAM_BUFFER stream_buffer = { .wb = NULL, };
512 if(unlikely(rrdhost_has_stream_sender_enabled(st->rrdhost)))
513 stream_buffer = stream_send_metrics_init(st, now.tv_sec);
514
515 spinlock_lock(&st->data_collection_lock);
516
517 if (pending_rrdset_next)
518 rrdset_timed_next(st, now, 0ULL);
519
520 netdata_log_debug(D_RRD_CALLS, "rrdset_done() for chart '%s'", rrdset_name(st));
521
522 RRDDIM *rd;
523
524 char
525 store_this_entry = 1, // boolean: 1 = store this entry, 0 = don't store this entry
526 first_entry = 0; // boolean: 1 = this is the first entry seen for this chart, 0 = all other entries
527
528 usec_t
529 last_collect_ut = 0, // the timestamp in microseconds, of the last collected value
530 now_collect_ut = 0, // the timestamp in microseconds, of this collected value (this is NOW)
531 last_stored_ut = 0, // the timestamp in microseconds, of the last stored entry in the db
532 next_store_ut = 0, // the timestamp in microseconds, of the next entry to store in the db
533 update_every_ut = st->update_every * USEC_PER_SEC; // st->update_every in microseconds
534
535 RRDSET_FLAGS rrdset_flags = rrdset_flag_check(st, ~0);
536 if(unlikely(rrdset_flags & RRDSET_FLAG_COLLECTION_FINISHED)) {
537 spinlock_unlock(&st->data_collection_lock);
538 return;
539 }
540
541 if (unlikely(rrdset_flags & RRDSET_FLAG_OBSOLETE)) {
542 netdata_log_error("Chart '%s' has the OBSOLETE flag set, but it is collected.", rrdset_id(st));
543 if(!spinlock_trylock(&st->destroy_lock))
544 fatal("RRDSET: chart '%s' of host '%s' is being collected while is being destroyed.", rrdset_id(st), rrdhost_hostname(st->rrdhost));
545
546 rrdset_isnot_obsolete___safe_from_collector_thread(st);
547 spinlock_unlock(&st->destroy_lock);
548 }
549
550 // check if the chart has a long time to be updated
551 if(unlikely(st->usec_since_last_update > MAX(st->db.entries, 60) * update_every_ut)) {
552 nd_log_daemon(NDLP_DEBUG, "host '%s', chart '%s': took too long to be updated (counter #%u, update #%u, %0.3" NETDATA_DOUBLE_MODIFIER
553 " secs). Resetting it.", rrdhost_hostname(st->rrdhost), rrdset_id(st), st->counter, st->counter_done,
554 (NETDATA_DOUBLE)st->usec_since_last_update / USEC_PER_SEC);
555 rrdset_collection_reset(st);
556 st->usec_since_last_update = update_every_ut;
557 store_this_entry = 0;
558 first_entry = 1;
559 }
560
561 rrdset_debug(st, "microseconds since last update: %"PRIu64"", st->usec_since_last_update);
562
563 // set last_collected_time
564 if(unlikely(!st->last_collected_time.tv_sec)) {
565 // it is the first entry
566 // set the last_collected_time to now
567 last_collect_ut = rrdset_init_last_collected_time(st, now) - update_every_ut;
568
569 // the first entry should not be stored
570 store_this_entry = 0;
571 first_entry = 1;
572 }
573 else {
574 // it is not the first entry
575 // calculate the proper last_collected_time, using usec_since_last_update
576 last_collect_ut = rrdset_update_last_collected_time(st);
577 }
578
579 // if this set has not been updated in the past
580 // we fake the last_update time to be = now - usec_since_last_update
581 if(unlikely(!st->last_updated.tv_sec)) {
582 // it has never been updated before
583 // set a fake last_updated, in the past using usec_since_last_update
584 rrdset_init_last_updated_time(st);
585
586 // the first entry should not be stored
587 store_this_entry = 0;
588 first_entry = 1;
589 }
590
591 // check if we will re-write the entire data set
592 if(unlikely(dt_usec(&st->last_collected_time, &st->last_updated) > st->db.entries * update_every_ut &&
593 st->rrd_memory_mode != RRD_DB_MODE_DBENGINE)) {
594 nd_log_daemon(NDLP_DEBUG, "'%s': too old data (last updated at %" PRId64 ".%" PRId64 ", last collected at %" PRId64 ".%" PRId64 "). "
595 "Resetting it. Will not store the next entry.",
596 rrdset_id(st),
597 (int64_t)st->last_updated.tv_sec,
598 (int64_t)st->last_updated.tv_usec,
599 (int64_t)st->last_collected_time.tv_sec,
600 (int64_t)st->last_collected_time.tv_usec);
601 rrdset_collection_reset(st);
602 rrdset_init_last_updated_time(st);
603
604 st->usec_since_last_update = update_every_ut;
605
606 // the first entry should not be stored
607 store_this_entry = 0;
608 first_entry = 1;
609 }
610
611 // these are the 3 variables that will help us in interpolation
612 // last_stored_ut = the last time we added a value to the storage
613 // now_collect_ut = the time the current value has been collected
614 // next_store_ut = the time of the next interpolation point
615 now_collect_ut = st->last_collected_time.tv_sec * USEC_PER_SEC + st->last_collected_time.tv_usec;
616 last_stored_ut = st->last_updated.tv_sec * USEC_PER_SEC + st->last_updated.tv_usec;
617 next_store_ut = (st->last_updated.tv_sec + st->update_every) * USEC_PER_SEC;
618
619 if(unlikely(!st->counter_done)) {
620 // set a fake last_updated to jump to current time
621 rrdset_init_last_updated_time(st);
622
623 last_stored_ut = st->last_updated.tv_sec * USEC_PER_SEC + st->last_updated.tv_usec;
624 next_store_ut = (st->last_updated.tv_sec + st->update_every) * USEC_PER_SEC;
625
626 if(unlikely(rrdset_flags & RRDSET_FLAG_STORE_FIRST)) {
627 store_this_entry = 1;
628 last_collect_ut = next_store_ut - update_every_ut;
629
630 rrdset_debug(st, "Fixed first entry.");
631 }
632 else {
633 store_this_entry = 0;
634
635 rrdset_debug(st, "Will not store the next entry.");
636 }
637 }
638
639 st->counter_done++;
640
641 if(stream_buffer.wb && !stream_buffer.v2)
642 stream_send_rrdset_metrics_v1(&stream_buffer, st);
643
644 size_t rda_slots = dictionary_entries(st->rrddim_root_index);
645 struct rda_item *rda_base = rrdset_thread_rda_get(&rda_slots);
646
647 size_t dim_id;
648 size_t dimensions = 0;
649 struct rda_item *rda = rda_base;
650 NETDATA_DOUBLE collected_total = 0.0;
651 NETDATA_DOUBLE last_collected_total = 0.0;
652 rrddim_foreach_read(rd, st) {
653 if(rd_dfe.counter >= rda_slots)
654 break;
655
656 rda = &rda_base[dimensions++];
657
658 // store the dimension in the array
659 rda->item = dictionary_acquired_item_dup(st->rrddim_root_index, rd_dfe.item);
660 rda->rd = dictionary_acquired_item_value(rda->item);
661 rda->reset_or_overflow = false;
662
663 // calculate totals
664 if(likely(rrddim_check_updated(rd))) {
665 // if the new is smaller than the old (an overflow, or reset), set the old equal to the new
666 // to reset the calculation (it will give zero as the calculation for this second)
667 if(unlikely(rd->algorithm == RRD_ALGORITHM_PCENT_OVER_DIFF_TOTAL && rrddim_last_collected_as_double(rd) > rrddim_collected_as_double(rd))) {
668 netdata_log_debug(D_RRD_STATS, "'%s' / '%s': RESET or OVERFLOW. Last collected value = " NETDATA_DOUBLE_FORMAT ", current = " NETDATA_DOUBLE_FORMAT
669 , rrdset_id(st)
670 , rrddim_name(rd)
671 , rrddim_last_collected_as_double(rd)
672 , rrddim_collected_as_double(rd)
673 );
674
675 if(!(rrddim_option_check(rd, RRDDIM_OPTION_DONT_DETECT_RESETS_OR_OVERFLOWS)))
676 rda->reset_or_overflow = true;
677
678 if(rrddim_is_float(rd))
679 rrddim_set_last_collected_float(rd, rrddim_collected_as_double(rd));
680 else
681 rrddim_set_last_collected_int(rd, rd->collector.collected.i.collected_value);
682 }
683
684 last_collected_total += rrddim_last_collected_as_double(rd);
685 collected_total += rrddim_collected_as_double(rd);
686
687 if(unlikely(rrddim_flag_check(rd, RRDDIM_FLAG_OBSOLETE))) {
688 netdata_log_error("Dimension %s in chart '%s' has the OBSOLETE flag set, but it is collected.", rrddim_name(rd), rrdset_id(st));
689 if(!spinlock_trylock(&rd->destroy_lock))
690 fatal("RRDSET: dimension '%s' of chart '%s' of host '%s' is being collected while is being destroyed.", rrddim_id(rd), rrdset_id(st), rrdhost_hostname(st->rrdhost));
691
692 rrddim_isnot_obsolete___safe_from_collector_thread(st, rd);
693 spinlock_unlock(&rd->destroy_lock);
694 }
695 }
696 }
697 rrddim_foreach_done(rd);
698 rda_slots = dimensions;
699
700 rrdset_debug(st, "last_collect_ut = %0.3" NETDATA_DOUBLE_MODIFIER " (last collection time)", (NETDATA_DOUBLE)last_collect_ut/USEC_PER_SEC);
701 rrdset_debug(st, "now_collect_ut = %0.3" NETDATA_DOUBLE_MODIFIER " (current collection time)", (NETDATA_DOUBLE)now_collect_ut/USEC_PER_SEC);
702 rrdset_debug(st, "last_stored_ut = %0.3" NETDATA_DOUBLE_MODIFIER " (last updated time)", (NETDATA_DOUBLE)last_stored_ut/USEC_PER_SEC);
703 rrdset_debug(st, "next_store_ut = %0.3" NETDATA_DOUBLE_MODIFIER " (next interpolation point)", (NETDATA_DOUBLE)next_store_ut/USEC_PER_SEC);
704
705 // process all dimensions to calculate their values
706 // based on the collected figures only
707 // at this stage we do not interpolate anything
708 for(dim_id = 0, rda = rda_base ; dim_id < rda_slots ; ++dim_id, ++rda) {
709 rd = rda->rd;
710 if(unlikely(!rd)) continue;
711
712 if(unlikely(!rrddim_check_updated(rd))) {
713 rd->collector.calculated_value = 0;
714 continue;
715 }
716
717 rrdset_debug(st, "%s: START "
718 " last_collected_value = " NETDATA_DOUBLE_FORMAT
719 " collected_value = " NETDATA_DOUBLE_FORMAT
720 " last_calculated_value = " NETDATA_DOUBLE_FORMAT
721 " calculated_value = " NETDATA_DOUBLE_FORMAT
722 , rrddim_name(rd)
723 , rrddim_last_collected_as_double(rd)
724 , rrddim_collected_as_double(rd)
725 , rd->collector.last_calculated_value
726 , rd->collector.calculated_value
727 );
728
729 switch(rd->algorithm) {
730 case RRD_ALGORITHM_ABSOLUTE:
731 rd->collector.calculated_value = rrddim_collected_as_double(rd)
732 * (NETDATA_DOUBLE)rd->multiplier
733 / (NETDATA_DOUBLE)rd->divisor;
734
735 rrdset_debug(st, "%s: CALC ABS/ABS-NO-IN " NETDATA_DOUBLE_FORMAT " = "
736 NETDATA_DOUBLE_FORMAT
737 " * " NETDATA_DOUBLE_FORMAT
738 " / " NETDATA_DOUBLE_FORMAT
739 , rrddim_name(rd)
740 , rd->collector.calculated_value
741 , rrddim_collected_as_double(rd)
742 , (NETDATA_DOUBLE)rd->multiplier
743 , (NETDATA_DOUBLE)rd->divisor
744 );
745 break;
746
747 case RRD_ALGORITHM_PCENT_OVER_ROW_TOTAL:
748 if(unlikely(!collected_total))
749 rd->collector.calculated_value = 0;
750 else
751 // the percentage of the current value
752 // over the total of all dimensions
753 rd->collector.calculated_value =
754 (NETDATA_DOUBLE)100
755 * rrddim_collected_as_double(rd)
756 / (NETDATA_DOUBLE)collected_total;
757
758 rrdset_debug(st, "%s: CALC PCENT-ROW " NETDATA_DOUBLE_FORMAT " = 100"
759 " * " NETDATA_DOUBLE_FORMAT
760 " / " NETDATA_DOUBLE_FORMAT
761 , rrddim_name(rd)
762 , rd->collector.calculated_value
763 , rrddim_collected_as_double(rd)
764 , (NETDATA_DOUBLE)collected_total
765 );
766 break;
767
768 case RRD_ALGORITHM_INCREMENTAL:
769 if(unlikely(rd->collector.counter <= 1)) {
770 rd->collector.calculated_value = 0;
771 continue;
772 }
773
774 if(rrddim_is_int(rd)) {
775 uint64_t last = (uint64_t)rd->collector.collected.i.last_collected_value;
776 uint64_t new = (uint64_t)rd->collector.collected.i.collected_value;
777 uint64_t max = (uint64_t)rd->collector.collected.i.collected_value_max;
778
779 // If the new is smaller than the old (overflow/reset), handle wrap
780 if(unlikely(last > new)) {
781 netdata_log_debug(D_RRD_STATS, "'%s' / '%s': RESET or OVERFLOW. Last collected value = " NETDATA_DOUBLE_FORMAT ", current = " NETDATA_DOUBLE_FORMAT
782 , rrdset_id(st)
783 , rrddim_name(rd)
784 , rrddim_last_collected_as_double(rd)
785 , rrddim_collected_as_double(rd));
786
787 if(!(rrddim_option_check(rd, RRDDIM_OPTION_DONT_DETECT_RESETS_OR_OVERFLOWS)))
788 rda->reset_or_overflow = true;
789
790 uint64_t cap = (max > 0x00000000FFFFFFFFULL) ? 0xFFFFFFFFFFFFFFFFULL : 0x00000000FFFFFFFFULL;
791 uint64_t delta = cap - last + new;
792 uint64_t max_acceptable_rate = (cap / 100) * MAX_INCREMENTAL_PERCENT_RATE;
793
794 if (delta < max_acceptable_rate) {
795 rd->collector.calculated_value +=
796 (NETDATA_DOUBLE) delta
797 * (NETDATA_DOUBLE) rd->multiplier
798 / (NETDATA_DOUBLE) rd->divisor;
799 } else {
800 rd->collector.calculated_value += 0;
801 }
802 }
803 else {
804 rd->collector.calculated_value +=
805 (NETDATA_DOUBLE)((int64_t)(rd->collector.collected.i.collected_value - rd->collector.collected.i.last_collected_value))
806 * (NETDATA_DOUBLE) rd->multiplier
807 / (NETDATA_DOUBLE) rd->divisor;
808 }
809 }
810 else {
811 NETDATA_DOUBLE last = rrddim_last_collected_as_double(rd);
812 NETDATA_DOUBLE cur = rrddim_collected_as_double(rd);
813 if(unlikely(cur < last)) {
814 if(!(rrddim_option_check(rd, RRDDIM_OPTION_DONT_DETECT_RESETS_OR_OVERFLOWS)))
815 rda->reset_or_overflow = true;
816 rd->collector.calculated_value += 0;
817 }
818 else {
819 rd->collector.calculated_value +=
820 (cur - last) * (NETDATA_DOUBLE) rd->multiplier / (NETDATA_DOUBLE) rd->divisor;
821 }
822 }
823
824 rrdset_debug(st, "%s: CALC INC PRE " NETDATA_DOUBLE_FORMAT " = ("
825 NETDATA_DOUBLE_FORMAT " - " NETDATA_DOUBLE_FORMAT
826 ")"
827 " * " NETDATA_DOUBLE_FORMAT
828 " / " NETDATA_DOUBLE_FORMAT
829 , rrddim_name(rd)
830 , rd->collector.calculated_value
831 , rrddim_collected_as_double(rd), rrddim_last_collected_as_double(rd)
832 , (NETDATA_DOUBLE)rd->multiplier
833 , (NETDATA_DOUBLE)rd->divisor
834 );
835 break;
836
837 case RRD_ALGORITHM_PCENT_OVER_DIFF_TOTAL:
838 if(unlikely(rd->collector.counter <= 1)) {
839 rd->collector.calculated_value = 0;
840 continue;
841 }
842
843 // the percentage of the current increment
844 // over the increment of all dimensions together
845 if(unlikely(collected_total == last_collected_total))
846 rd->collector.calculated_value = 0;
847 else
848 rd->collector.calculated_value =
849 (NETDATA_DOUBLE)100
850 * (NETDATA_DOUBLE)(rrddim_collected_as_double(rd) - rrddim_last_collected_as_double(rd))
851 / (NETDATA_DOUBLE)(collected_total - last_collected_total);
852
853 rrdset_debug(st, "%s: CALC PCENT-DIFF " NETDATA_DOUBLE_FORMAT " = 100"
854 " * (" NETDATA_DOUBLE_FORMAT " - " NETDATA_DOUBLE_FORMAT ")"
855 " / (" NETDATA_DOUBLE_FORMAT " - " NETDATA_DOUBLE_FORMAT ")"
856 , rrddim_name(rd)
857 , rd->collector.calculated_value
858 , rrddim_collected_as_double(rd), rrddim_last_collected_as_double(rd)
859 , (NETDATA_DOUBLE)collected_total, (NETDATA_DOUBLE)last_collected_total
860 );
861 break;
862
863 default:
864 // make the default zero, to make sure
865 // it gets noticed when we add new types
866 rd->collector.calculated_value = 0;
867
868 rrdset_debug(st, "%s: CALC " NETDATA_DOUBLE_FORMAT " = 0"
869 , rrddim_name(rd)
870 , rd->collector.calculated_value
871 );
872 break;
873 }
874
875 rrdset_debug(st, "%s: PHASE2 "
876 " last_collected_value = " NETDATA_DOUBLE_FORMAT
877 " collected_value = " NETDATA_DOUBLE_FORMAT
878 " last_calculated_value = " NETDATA_DOUBLE_FORMAT
879 " calculated_value = " NETDATA_DOUBLE_FORMAT
880 , rrddim_name(rd)
881 , rrddim_last_collected_as_double(rd)
882 , rrddim_collected_as_double(rd)
883 , rd->collector.last_calculated_value
884 , rd->collector.calculated_value
885 );
886 }
887
888 // at this point we have all the calculated values ready
889 // it is now time to interpolate values on a second boundary
890
891 // #ifdef NETDATA_INTERNAL_CHECKS
892 // if(unlikely(now_collect_ut < next_store_ut && st->counter_done > 1)) {
893 // // this is collected in the same interpolation point
894 // rrdset_debug(st, "THIS IS IN THE SAME INTERPOLATION POINT");
895 // netdata_log_info("INTERNAL CHECK: host '%s', chart '%s' collection %zu is in the same interpolation point: short by %llu microseconds", st->rrdhost->hostname, rrdset_name(st), st->counter_done, next_store_ut - now_collect_ut);
896 // }
897 // #endif
898
899 rrdset_done_interpolate(
900 &stream_buffer
901 , st
902 , rda_base
903 , rda_slots
904 , update_every_ut
905 , last_stored_ut
906 , next_store_ut
907 , last_collect_ut
908 , now_collect_ut
909 , store_this_entry
910 );
911
912 for(dim_id = 0, rda = rda_base ; dim_id < rda_slots ; ++dim_id, ++rda) {
913 rd = rda->rd;
914 if(unlikely(!rd)) continue;
915
916 if(unlikely(!rrddim_check_updated(rd)))
917 continue;
918
919 rrdset_debug(st, "%s: setting last_collected_value (old: " NETDATA_DOUBLE_FORMAT ") to last_collected_value (new: " NETDATA_DOUBLE_FORMAT ")", rrddim_name(rd), rrddim_last_collected_as_double(rd), rrddim_collected_as_double(rd));
920
921 if(rrddim_is_float(rd))
922 rrddim_set_last_collected_float(rd, rrddim_collected_as_double(rd));
923 else
924 rrddim_set_last_collected_int(rd, rd->collector.collected.i.collected_value);
925
926 switch(rd->algorithm) {
927 case RRD_ALGORITHM_INCREMENTAL:
928 if(unlikely(!first_entry)) {
929 rrdset_debug(st, "%s: setting last_calculated_value (old: " NETDATA_DOUBLE_FORMAT ") to "
930 "last_calculated_value (new: " NETDATA_DOUBLE_FORMAT ")"
931 , rrddim_name(rd)
932 , rd->collector.last_calculated_value + rd->collector.calculated_value
933 , rd->collector.calculated_value);
934
935 rd->collector.last_calculated_value += rd->collector.calculated_value;
936 }
937 else {
938 rrdset_debug(st, "THIS IS THE FIRST POINT");
939 }
940 break;
941
942 case RRD_ALGORITHM_ABSOLUTE:
943 case RRD_ALGORITHM_PCENT_OVER_ROW_TOTAL:
944 case RRD_ALGORITHM_PCENT_OVER_DIFF_TOTAL:
945 rrdset_debug(st, "%s: setting last_calculated_value (old: " NETDATA_DOUBLE_FORMAT ") to "
946 "last_calculated_value (new: " NETDATA_DOUBLE_FORMAT ")"
947 , rrddim_name(rd)
948 , rd->collector.last_calculated_value
949 , rd->collector.calculated_value);
950
951 rd->collector.last_calculated_value = rd->collector.calculated_value;
952 break;
953 }
954
955 rd->collector.calculated_value = 0;
956 if(rrddim_is_float(rd))
957 rrddim_set_collected_float(rd, 0.0);
958 else
959 rrddim_set_collected_int(rd, 0);
960 rrddim_clear_updated(rd);
961
962 rrdset_debug(st, "%s: END "
963 " last_collected_value = " NETDATA_DOUBLE_FORMAT
964 " collected_value = " NETDATA_DOUBLE_FORMAT
965 " last_calculated_value = " NETDATA_DOUBLE_FORMAT
966 " calculated_value = " NETDATA_DOUBLE_FORMAT
967 , rrddim_name(rd)
968 , rrddim_last_collected_as_double(rd)
969 , rrddim_collected_as_double(rd)
970 , rd->collector.last_calculated_value
971 , rd->collector.calculated_value
972 );
973 }
974
975 spinlock_unlock(&st->data_collection_lock);
976 stream_send_rrdset_metrics_finished(&stream_buffer, st);
977
978 // ALL DONE ABOUT THE DATA UPDATE
979 // --------------------------------------------------------------------
980
981 for(dim_id = 0, rda = rda_base; dim_id < rda_slots ; ++dim_id, ++rda) {
982 rd = rda->rd;
983 if(unlikely(!rd)) continue;
984
985 dictionary_acquired_item_release(st->rrddim_root_index, rda->item);
986 rda->item = NULL;
987 rda->rd = NULL;
988 }
989
990 rrdcontext_collected_rrdset(st);
991
992 store_metric_collection_completed();
993 }