@cryptotaxi247 / netdata-1 / commits / 9b83b2b2f

journal estimations (#16445)

* preliminary work to enable estimations * working estimations * clear estimated counter when switching files * use the actual query timeframe for estimations * properly estimate the remaining lines in the files * fix the query duration according to file limits * require at least 2 seconds of data per file * estimation starts at 5% of each file * work on estimation * distribute estimations accurately on the histogram * stop query permaturely if it is going to timeout * more aggressive estimation of next file duration * added missing dimension names in histogram.chart.summary.dimension; reordered estimated and unsampled; added histogram.chart.view.dimensions.colors to estimated and unsampled; added message * include message always * message is sent on non-data-only queries; empty message object when there is nothing to tell * use syslog priorities * use standard netdata log priorities

Costa Tsaousis committed Nov 22, 2023 at 15:50 UTC 9b83b2b2f5fd1e4e2d2ec8224aefba9ee19c4b87
7 files changed +529 -149
collectors/systemd-journal.plugin/systemd-internals.h
+1
@@ -83,6 +83,7 @@ int journal_file_dict_items_backward_compar(const void *a, const void *b);
83 int journal_file_dict_items_forward_compar(const void *a, const void *b);
84 void buffer_json_journal_versions(BUFFER *wb);
85 void available_journal_file_sources_to_json_array(BUFFER *wb);
86 +bool journal_files_completed_once(void);
87
88 FACET_ROW_SEVERITY syslog_priority_to_facet_severity(FACETS *facets, FACET_ROW *row, void *data);
89
collectors/systemd-journal.plugin/systemd-journal-files.c
+6
@@ -391,6 +391,11 @@ void journal_directory_scan(const char *dirname, int depth, usec_t last_scan_ut)
391 closedir(dir);
392 }
393
394 +static size_t journal_files_scans = 0;
395 +bool journal_files_completed_once(void) {
396 + return journal_files_scans > 0;
397 +}
398 +
399 void journal_files_registry_update(void) {
400 static SPINLOCK spinlock = NETDATA_SPINLOCK_INITIALIZER;
401
@@ -411,6 +416,7 @@ void journal_files_registry_update(void) {
416 }
417 dfe_done(jf);
418
419 + journal_files_scans++;
420 spinlock_unlock(&spinlock);
421 }
422 }
collectors/systemd-journal.plugin/systemd-journal.c
+290 -58
@@ -211,11 +211,18 @@ typedef struct function_query_status {
211 const char *query;
212 const char *histogram;
213
214 + struct {
215 + usec_t start_ut; // the starting time of the query - we start from this
216 + usec_t stop_ut; // the ending time of the query - we stop at this
217 + usec_t first_msg_ut;
218 + } query_file;
219 +
220 struct {
221 uint32_t enable_after_samples;
222 uint32_t slots;
223 uint32_t sampled;
224 uint32_t unsampled;
225 + uint32_t estimated;
226 } samples;
227
228 struct {
@@ -225,6 +232,7 @@ typedef struct function_query_status {
232 uint32_t recalibrate;
233 uint32_t sampled;
234 uint32_t unsampled;
235 + uint32_t estimated;
236 } samples_per_file;
237
238 struct {
@@ -328,10 +336,14 @@ static void sampling_query_init(FUNCTION_QUERY_STATUS *fqs, FACETS *facets) {
336 // the minimum number of rows to enable sampling
337 fqs->samples.enable_after_samples = fqs->sampling / 2;
338
339 + size_t files_matched = fqs->files_matched;
340 + if(!files_matched)
341 + files_matched = 1;
342 +
343 // the minimum number of rows per file to enable sampling
332 - fqs->samples_per_file.enable_after_samples = (fqs->sampling / 4) / fqs->files_matched;
333 - if(fqs->samples_per_file.enable_after_samples < fqs->entries * 2)
334 - fqs->samples_per_file.enable_after_samples = fqs->entries * 2;
344 + fqs->samples_per_file.enable_after_samples = (fqs->sampling / 4) / files_matched;
345 + if(fqs->samples_per_file.enable_after_samples < fqs->entries)
346 + fqs->samples_per_file.enable_after_samples = fqs->entries;
347
348 // the minimum number of rows per time slot to enable sampling
349 fqs->samples_per_time_slot.enable_after_samples = (fqs->sampling / 4) / fqs->samples.slots;
@@ -342,57 +354,156 @@ static void sampling_query_init(FUNCTION_QUERY_STATUS *fqs, FACETS *facets) {
354 static void sampling_file_init(FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf __maybe_unused) {
355 fqs->samples_per_file.sampled = 0;
356 fqs->samples_per_file.unsampled = 0;
357 + fqs->samples_per_file.estimated = 0;
358 fqs->samples_per_file.every = 0;
359 fqs->samples_per_file.skipped = 0;
360 fqs->samples_per_file.recalibrate = 0;
361 }
362
350 -static void sampling_decide_file_sampling_every(FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf, FACETS_ANCHOR_DIRECTION direction, usec_t msg_ut) {
363 +static size_t sampling_file_lines_scanned(FUNCTION_QUERY_STATUS *fqs) {
364 size_t sampled = fqs->samples_per_file.sampled + fqs->samples_per_file.unsampled;
365 + if(!sampled) sampled = 1;
366 + return sampled;
367 +}
368
353 - if(!sampled)
354 - sampled = 1;
369 +static double sampling_file_query_overlapping_timeframe_ut(FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf, FACETS_ANCHOR_DIRECTION direction,
370 + usec_t msg_ut, usec_t *after_ut, usec_t *before_ut) {
371 + // find the overlap of the query and file timeframes
372 + // taking into account the first message we encountered
373 +
374 + double progress;
375 + usec_t oldest_ut, newest_ut, elapsed_ut;
376 + if(direction == FACETS_ANCHOR_DIRECTION_FORWARD) {
377 + // the first message we know (oldest)
378 + oldest_ut = fqs->query_file.first_msg_ut ? fqs->query_file.first_msg_ut : jf->msg_first_ut;
379 + if(!oldest_ut) oldest_ut = fqs->query_file.start_ut;
380 +
381 + if(jf->msg_last_ut)
382 + newest_ut = MIN(fqs->query_file.stop_ut, jf->msg_last_ut);
383 + else if(jf->file_last_modified_ut)
384 + newest_ut = MIN(fqs->query_file.stop_ut, jf->file_last_modified_ut);
385 + else
386 + newest_ut = fqs->query_file.stop_ut;
387 +
388 + if(msg_ut < oldest_ut)
389 + oldest_ut = msg_ut - 1;
390 +
391 + elapsed_ut = msg_ut - oldest_ut;
392 + }
393 + else /* BACKWARD */ {
394 + // the latest message we know (newest)
395 + newest_ut = fqs->query_file.first_msg_ut ? fqs->query_file.first_msg_ut : jf->msg_last_ut;
396 + if(!newest_ut) newest_ut = fqs->query_file.start_ut;
397 +
398 + if(jf->msg_first_ut)
399 + oldest_ut = MAX(fqs->query_file.stop_ut, jf->msg_first_ut);
400 + else
401 + oldest_ut = fqs->query_file.stop_ut;
402
356 - // find the common duration
357 - usec_t after_ut = fqs->samples_per_time_slot.start_ut < jf->msg_first_ut ? jf->msg_first_ut : fqs->samples_per_time_slot.start_ut;
358 - usec_t before_ut = fqs->samples_per_time_slot.end_ut > jf->msg_last_ut ? jf->msg_last_ut : fqs->samples_per_time_slot.end_ut;
403 + if(newest_ut < msg_ut)
404 + newest_ut = msg_ut + 1;
405
360 - if(after_ut > before_ut) {
361 - usec_t t = after_ut;
362 - after_ut = before_ut;
363 - before_ut = t;
406 + elapsed_ut = newest_ut - msg_ut;
407 }
408
366 - if(after_ut == before_ut)
367 - after_ut = before_ut - 1;
409 + usec_t total_ut = newest_ut - oldest_ut;
410 + progress = (double)elapsed_ut / (double)total_ut;
411
412 + *after_ut = oldest_ut;
413 + *before_ut = newest_ut;
414 +
415 + return progress;
416 +}
417 +
418 +static usec_t sampling_file_remaining_time_ut(FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf, FACETS_ANCHOR_DIRECTION direction,
419 + usec_t msg_ut, usec_t *total_time_ut, usec_t *remaining_start_ut, usec_t *remaining_end_ut) {
420 +
421 + usec_t after_ut, before_ut;
422 + sampling_file_query_overlapping_timeframe_ut(fqs, jf, direction, msg_ut, &after_ut, &before_ut);
423 +
424 + // since we have a timestamp in msg_ut
425 + // this timestamp can extend the overlap
426 if(msg_ut <= after_ut)
370 - msg_ut = after_ut + 1;
427 + after_ut = msg_ut - 1;
428
429 if(msg_ut >= before_ut)
373 - msg_ut = before_ut - 1;
430 + before_ut = msg_ut + 1;
431 +
432 + // return the remaining duration
433 + usec_t remaining_from_ut, remaining_to_ut;
434 + if(direction == FACETS_ANCHOR_DIRECTION_FORWARD) {
435 + remaining_from_ut = msg_ut;
436 + remaining_to_ut = before_ut;
437 + }
438 + else {
439 + remaining_from_ut = after_ut;
440 + remaining_to_ut = msg_ut;
441 + }
442
375 - size_t expected_lines;
443 + usec_t remaining_ut = remaining_to_ut - remaining_from_ut;
444
377 - if(direction == FACETS_ANCHOR_DIRECTION_FORWARD)
378 - expected_lines = sampled * (before_ut - after_ut) / (msg_ut - after_ut);
379 - else
380 - expected_lines = sampled * (before_ut - after_ut) / (before_ut - msg_ut);
445 + if(total_time_ut)
446 + *total_time_ut = (before_ut > after_ut) ? before_ut - after_ut : 1;
447 +
448 + if(remaining_start_ut)
449 + *remaining_start_ut = remaining_from_ut;
450 +
451 + if(remaining_end_ut)
452 + *remaining_end_ut = remaining_to_ut;
453
382 - if(expected_lines < 1)
454 + return remaining_ut;
455 +}
456 +
457 +static size_t sampling_file_estimate_remaining_lines(FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf, FACETS_ANCHOR_DIRECTION direction, usec_t msg_ut) {
458 + size_t scanned_lines = sampling_file_lines_scanned(fqs);
459 +
460 + usec_t total_time_ut;
461 + usec_t remaining_time_ut = sampling_file_remaining_time_ut(fqs, jf, direction, msg_ut, &total_time_ut, NULL, NULL);
462 +
463 + if (total_time_ut == 0)
464 + total_time_ut = 1;
465 +
466 + // Calculate the proportion of time covered
467 + double time_proportion = (double)(total_time_ut - remaining_time_ut) / (double)total_time_ut;
468 +
469 + if (time_proportion == 0 || !isfinite(time_proportion))
470 + time_proportion = 1.0;
471 +
472 + // Estimate the total number of lines in the file
473 + size_t total_lines_estimated = (size_t)((double)scanned_lines / time_proportion);
474 +
475 + // Calculate the estimated number of remaining lines
476 + size_t expected_lines = total_lines_estimated - scanned_lines;
477 +
478 + if (expected_lines < 1)
479 expected_lines = 1;
480
385 - size_t wanted_samples = (fqs->sampling / 2) / fqs->files_matched;
481 + return expected_lines;
482 +}
483
387 - fqs->samples_per_file.every = expected_lines / wanted_samples;
484 +static void sampling_decide_file_sampling_every(FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf, FACETS_ANCHOR_DIRECTION direction, usec_t msg_ut) {
485 + size_t files_matched = fqs->files_matched;
486 + if(!files_matched) files_matched = 1;
487 +
488 + size_t remaining_lines = sampling_file_estimate_remaining_lines(fqs, jf, direction, msg_ut);
489 + size_t wanted_samples = (fqs->sampling / 2) / files_matched;
490 + if(!wanted_samples) wanted_samples = 1;
491 +
492 + fqs->samples_per_file.every = remaining_lines / wanted_samples;
493
494 if(fqs->samples_per_file.every < 1)
495 fqs->samples_per_file.every = 1;
496 }
497
393 -static inline bool is_row_in_sample(FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf, usec_t msg_ut, FACETS_ANCHOR_DIRECTION direction, bool candidate_to_keep) {
394 - if(!fqs->sampling)
395 - return true;
498 +typedef enum {
499 + SAMPLING_STOP_AND_ESTIMATE = -1,
500 + SAMPLING_FULL = 0,
501 + SAMPLING_SKIP_FIELDS = 1,
502 +} sampling_t;
503 +
504 +static inline sampling_t is_row_in_sample(FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf, usec_t msg_ut, FACETS_ANCHOR_DIRECTION direction, bool candidate_to_keep) {
505 + if(!fqs->sampling || candidate_to_keep)
506 + return SAMPLING_FULL;
507
508 if(unlikely(msg_ut < fqs->samples_per_time_slot.start_ut))
509 msg_ut = fqs->samples_per_time_slot.start_ut;
@@ -403,7 +514,7 @@ static inline bool is_row_in_sample(FUNCTION_QUERY_STATUS *fqs, struct journal_f
514 if(slot >= fqs->samples.slots)
515 slot = fqs->samples.slots - 1;
516
406 - bool should_sample = candidate_to_keep;
517 + bool should_sample = false;
518
519 if(fqs->samples.sampled < fqs->samples.enable_after_samples ||
520 fqs->samples_per_file.sampled < fqs->samples_per_file.enable_after_samples ||
@@ -426,20 +537,40 @@ static inline bool is_row_in_sample(FUNCTION_QUERY_STATUS *fqs, struct journal_f
537 fqs->samples_per_file.skipped++;
538 }
539
429 - fqs->samples_per_file.recalibrate++;
430 -
540 if(should_sample) {
541 fqs->samples.sampled++;
542 fqs->samples_per_file.sampled++;
543 fqs->samples_per_time_slot.sampled[slot]++;
544 +
545 + return SAMPLING_FULL;
546 }
436 - else {
437 - fqs->samples.unsampled++;
438 - fqs->samples_per_file.unsampled++;
439 - fqs->samples_per_time_slot.unsampled[slot]++;
547 +
548 + fqs->samples_per_file.recalibrate++;
549 +
550 + fqs->samples.unsampled++;
551 + fqs->samples_per_file.unsampled++;
552 + fqs->samples_per_time_slot.unsampled[slot]++;
553 +
554 + if(fqs->samples_per_file.unsampled > fqs->samples_per_file.sampled) {
555 + usec_t after_ut, before_ut;
556 +
557 + double progress = sampling_file_query_overlapping_timeframe_ut(
558 + fqs, jf, direction, msg_ut, &after_ut, &before_ut);
559 +
560 + if(progress > 0.01)
561 + return SAMPLING_STOP_AND_ESTIMATE;
562 }
563
442 - return should_sample;
564 + return SAMPLING_SKIP_FIELDS;
565 +}
566 +
567 +static void sampling_update_file_estimates(FACETS *facets, FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf, usec_t msg_ut, FACETS_ANCHOR_DIRECTION direction) {
568 + usec_t total_time_ut, remaining_start_ut, remaining_end_ut;
569 + sampling_file_remaining_time_ut(fqs, jf, direction, msg_ut, &total_time_ut, &remaining_start_ut, &remaining_end_ut);
570 + size_t remaining_lines = sampling_file_estimate_remaining_lines(fqs, jf, direction, msg_ut);
571 + facets_update_estimations(facets, remaining_start_ut, remaining_end_ut, remaining_lines);
572 + fqs->samples.estimated += remaining_lines;
573 + fqs->samples_per_file.estimated += remaining_lines;
574 }
575
576 // ----------------------------------------------------------------------------
@@ -539,11 +670,15 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_backward(
670 usec_t stop_ut = (fqs->data_only && fqs->anchor.stop_ut) ? fqs->anchor.stop_ut : fqs->after_ut;
671 bool stop_when_full = (fqs->data_only && !fqs->anchor.stop_ut);
672
673 + fqs->query_file.start_ut = start_ut;
674 + fqs->query_file.stop_ut = stop_ut;
675 +
676 if(!netdata_systemd_journal_seek_to(j, start_ut))
677 return ND_SD_JOURNAL_FAILED_TO_SEEK;
678
679 size_t errors_no_timestamp = 0;
546 - usec_t earliest_msg_ut = 0;
680 + usec_t latest_msg_ut = 0; // the biggest timestamp we have seen so far
681 + usec_t first_msg_ut = 0; // the first message we got from the db
682 size_t row_counter = 0, last_row_counter = 0, rows_useful = 0;
683 size_t bytes = 0, last_bytes = 0;
684
@@ -560,20 +695,25 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_backward(
695 continue;
696 }
697
563 - if(unlikely(msg_ut > earliest_msg_ut))
564 - earliest_msg_ut = msg_ut;
565 -
698 if (unlikely(msg_ut > start_ut))
699 continue;
700
701 if (unlikely(msg_ut < stop_ut))
702 break;
703
572 - bool to_sample = is_row_in_sample(fqs, jf, msg_ut,
704 + if(unlikely(msg_ut > latest_msg_ut))
705 + latest_msg_ut = msg_ut;
706 +
707 + if(unlikely(!first_msg_ut)) {
708 + first_msg_ut = msg_ut;
709 + fqs->query_file.first_msg_ut = msg_ut;
710 + }
711 +
712 + sampling_t sample = is_row_in_sample(fqs, jf, msg_ut,
713 FACETS_ANCHOR_DIRECTION_BACKWARD,
714 facets_row_candidate_to_keep(facets, msg_ut));
715
576 - if(to_sample) {
716 + if(sample == SAMPLING_FULL) {
717 bytes += netdata_systemd_journal_process_row(j, facets, jf, &msg_ut);
718
719 // make sure each line gets a unique timestamp
@@ -605,8 +745,12 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_backward(
745 status = check_stop(fqs->cancelled, &fqs->stop_monotonic_ut);
746 }
747 }
608 - else
748 + else if(sample == SAMPLING_SKIP_FIELDS)
749 facets_row_finished_unsampled(facets, msg_ut);
750 + else {
751 + sampling_update_file_estimates(facets, fqs, jf, msg_ut, FACETS_ANCHOR_DIRECTION_BACKWARD);
752 + break;
753 + }
754 }
755
756 FUNCTION_PROGRESS_UPDATE_ROWS(fqs->rows_read, row_counter - last_row_counter);
@@ -617,8 +761,8 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_backward(
761 if(errors_no_timestamp)
762 netdata_log_error("SYSTEMD-JOURNAL: %zu lines did not have timestamps", errors_no_timestamp);
763
620 - if(earliest_msg_ut > fqs->last_modified)
621 - fqs->last_modified = earliest_msg_ut;
764 + if(latest_msg_ut > fqs->last_modified)
765 + fqs->last_modified = latest_msg_ut;
766
767 return status;
768 }
@@ -633,11 +777,15 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_forward(
777 usec_t stop_ut = ((fqs->data_only && fqs->anchor.stop_ut) ? fqs->anchor.stop_ut : fqs->before_ut) + anchor_delta;
778 bool stop_when_full = (fqs->data_only && !fqs->anchor.stop_ut);
779
780 + fqs->query_file.start_ut = start_ut;
781 + fqs->query_file.stop_ut = stop_ut;
782 +
783 if(!netdata_systemd_journal_seek_to(j, start_ut))
784 return ND_SD_JOURNAL_FAILED_TO_SEEK;
785
786 size_t errors_no_timestamp = 0;
640 - usec_t earliest_msg_ut = 0;
787 + usec_t latest_msg_ut = 0; // the biggest timestamp we have seen so far
788 + usec_t first_msg_ut = 0; // the first message we got from the db
789 size_t row_counter = 0, last_row_counter = 0, rows_useful = 0;
790 size_t bytes = 0, last_bytes = 0;
791
@@ -654,20 +802,25 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_forward(
802 continue;
803 }
804
657 - if(likely(msg_ut > earliest_msg_ut))
658 - earliest_msg_ut = msg_ut;
659 -
805 if (unlikely(msg_ut < start_ut))
806 continue;
807
808 if (unlikely(msg_ut > stop_ut))
809 break;
810
666 - bool to_sample = is_row_in_sample(fqs, jf, msg_ut,
811 + if(likely(msg_ut > latest_msg_ut))
812 + latest_msg_ut = msg_ut;
813 +
814 + if(unlikely(!first_msg_ut)) {
815 + first_msg_ut = msg_ut;
816 + fqs->query_file.first_msg_ut = msg_ut;
817 + }
818 +
819 + sampling_t sample = is_row_in_sample(fqs, jf, msg_ut,
820 FACETS_ANCHOR_DIRECTION_FORWARD,
821 facets_row_candidate_to_keep(facets, msg_ut));
822
670 - if(to_sample) {
823 + if(sample == SAMPLING_FULL) {
824 bytes += netdata_systemd_journal_process_row(j, facets, jf, &msg_ut);
825
826 // make sure each line gets a unique timestamp
@@ -699,8 +852,12 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_forward(
852 status = check_stop(fqs->cancelled, &fqs->stop_monotonic_ut);
853 }
854 }
702 - else
855 + else if(sample == SAMPLING_SKIP_FIELDS)
856 facets_row_finished_unsampled(facets, msg_ut);
857 + else {
858 + sampling_update_file_estimates(facets, fqs, jf, msg_ut, FACETS_ANCHOR_DIRECTION_FORWARD);
859 + break;
860 + }
861 }
862
863 FUNCTION_PROGRESS_UPDATE_ROWS(fqs->rows_read, row_counter - last_row_counter);
@@ -711,8 +868,8 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_forward(
868 if(errors_no_timestamp)
869 netdata_log_error("SYSTEMD-JOURNAL: %zu lines did not have timestamps", errors_no_timestamp);
870
714 - if(earliest_msg_ut > fqs->last_modified)
715 - fqs->last_modified = earliest_msg_ut;
871 + if(latest_msg_ut > fqs->last_modified)
872 + fqs->last_modified = latest_msg_ut;
873
874 return status;
875 }
@@ -919,8 +1076,10 @@ static int netdata_systemd_journal_query(BUFFER *wb, FACETS *facets, FUNCTION_QU
1076 }
1077
1078 bool partial = false;
922 - usec_t started_ut;
923 - usec_t ended_ut = now_monotonic_usec();
1079 + usec_t query_started_ut = now_monotonic_usec();
1080 + usec_t started_ut = query_started_ut;
1081 + usec_t ended_ut = started_ut;
1082 + usec_t duration_ut = 0, max_duration_ut = 0;
1083
1084 sampling_query_init(fqs, facets);
1085
@@ -932,6 +1091,17 @@ static int netdata_systemd_journal_query(BUFFER *wb, FACETS *facets, FUNCTION_QU
1091 if(!jf_is_mine(jf, fqs))
1092 continue;
1093
1094 + started_ut = ended_ut;
1095 +
1096 + // do not even try to do the query if we expect it to pass the timeout
1097 + if(ended_ut > (query_started_ut + (fqs->stop_monotonic_ut - query_started_ut) * 3 / 4) &&
1098 + ended_ut + max_duration_ut * 2 >= fqs->stop_monotonic_ut) {
1099 +
1100 + partial = true;
1101 + status = ND_SD_JOURNAL_TIMED_OUT;
1102 + break;
1103 + }
1104 +
1105 fqs->file_working++;
1106 // fqs->cached_count = 0;
1107
@@ -953,9 +1123,11 @@ static int netdata_systemd_journal_query(BUFFER *wb, FACETS *facets, FUNCTION_QU
1123 fs_calls = fstat_thread_calls - fs_calls;
1124 fs_cached = fstat_thread_cached_responses - fs_cached;
1125
956 - started_ut = ended_ut;
1126 ended_ut = now_monotonic_usec();
958 - usec_t duration_ut = ended_ut - started_ut;
1127 + duration_ut = ended_ut - started_ut;
1128 +
1129 + if(duration_ut > max_duration_ut)
1130 + max_duration_ut = duration_ut;
1131
1132 buffer_json_add_array_item_object(wb); // journal file
1133 {
@@ -984,6 +1156,7 @@ static int netdata_systemd_journal_query(BUFFER *wb, FACETS *facets, FUNCTION_QU
1156 {
1157 buffer_json_member_add_uint64(wb, "sampled", fqs->samples_per_file.sampled);
1158 buffer_json_member_add_uint64(wb, "unsampled", fqs->samples_per_file.unsampled);
1159 + buffer_json_member_add_uint64(wb, "estimated", fqs->samples_per_file.estimated);
1160 }
1161 buffer_json_object_close(wb); // _sampling
1162 }
@@ -1056,6 +1229,64 @@ static int netdata_systemd_journal_query(BUFFER *wb, FACETS *facets, FUNCTION_QU
1229 buffer_json_member_add_boolean(wb, "partial", partial);
1230 buffer_json_member_add_string(wb, "type", "table");
1231
1232 + // build a message for the query
1233 + if(!fqs->data_only) {
1234 + CLEAN_BUFFER *msg = buffer_create(0, NULL);
1235 + CLEAN_BUFFER *msg_description = buffer_create(0, NULL);
1236 + ND_LOG_FIELD_PRIORITY msg_priority = NDLP_INFO;
1237 +
1238 + if(!journal_files_completed_once()) {
1239 + buffer_strcat(msg, "Journals are still being scanned. ");
1240 + buffer_strcat(msg_description
1241 + , "LIBRARY SCAN: The journal files are still being scanned, you are probably viewing incomplete data. ");
1242 + msg_priority = NDLP_WARNING;
1243 + }
1244 +
1245 + if(partial) {
1246 + buffer_strcat(msg, "Query timed-out, incomplete data. ");
1247 + buffer_strcat(msg_description
1248 + , "QUERY TIMEOUT: The query timed out and may not include all the data of the selected window. ");
1249 + msg_priority = NDLP_WARNING;
1250 + }
1251 +
1252 + if(fqs->samples.estimated || fqs->samples.unsampled) {
1253 + double percent = (double) (fqs->samples.sampled * 100.0 /
1254 + (fqs->samples.estimated + fqs->samples.unsampled + fqs->samples.sampled));
1255 + buffer_sprintf(msg, "%.2f%% real data", percent);
1256 + buffer_sprintf(msg_description, "ACTUAL DATA: The filters counters reflect %0.2f%% of the data. ", percent);
1257 + msg_priority = MIN(msg_priority, NDLP_NOTICE);
1258 + }
1259 +
1260 + if(fqs->samples.unsampled) {
1261 + double percent = (double) (fqs->samples.unsampled * 100.0 /
1262 + (fqs->samples.estimated + fqs->samples.unsampled + fqs->samples.sampled));
1263 + buffer_sprintf(msg, ", %.2f%% unsampled", percent);
1264 + buffer_sprintf(msg_description
1265 + , "UNSAMPLED DATA: %0.2f%% of the events exist and have been counted, but their values have not been evaluated, so they are not included in the filters counters. "
1266 + , percent);
1267 + msg_priority = MIN(msg_priority, NDLP_NOTICE);
1268 + }
1269 +
1270 + if(fqs->samples.estimated) {
1271 + double percent = (double) (fqs->samples.estimated * 100.0 /
1272 + (fqs->samples.estimated + fqs->samples.unsampled + fqs->samples.sampled));
1273 + buffer_sprintf(msg, ", %.2f%% estimated", percent);
1274 + buffer_sprintf(msg_description
1275 + , "ESTIMATED DATA: The query selected a large amount of data, so to avoid delaying too much, the presented data are estimated by %0.2f%%. "
1276 + , percent);
1277 + msg_priority = MIN(msg_priority, NDLP_NOTICE);
1278 + }
1279 +
1280 + buffer_json_member_add_object(wb, "message");
1281 + if(buffer_tostring(msg)) {
1282 + buffer_json_member_add_string(wb, "title", buffer_tostring(msg));
1283 + buffer_json_member_add_string(wb, "description", buffer_tostring(msg_description));
1284 + buffer_json_member_add_string(wb, "status", nd_log_id2priority(msg_priority));
1285 + }
1286 + // else send an empty object if there is nothing to tell
1287 + buffer_json_object_close(wb); // message
1288 + }
1289 +
1290 if(!fqs->data_only) {
1291 buffer_json_member_add_time_t(wb, "update_every", 1);
1292 buffer_json_member_add_string(wb, "help", SYSTEMD_JOURNAL_FUNCTION_DESCRIPTION);
@@ -1081,6 +1312,7 @@ static int netdata_systemd_journal_query(BUFFER *wb, FACETS *facets, FUNCTION_QU
1312 {
1313 buffer_json_member_add_uint64(wb, "sampled", fqs->samples.sampled);
1314 buffer_json_member_add_uint64(wb, "unsampled", fqs->samples.unsampled);
1315 + buffer_json_member_add_uint64(wb, "estimated", fqs->samples.estimated);
1316 }
1317 buffer_json_object_close(wb); // _sampling
1318 }
libnetdata/facets/facets.c
+228 -90
@@ -36,7 +36,8 @@ static const uint8_t id_encoding_characters_reverse[256] = {
36 #define FACETS_HASH XXH64_hash_t
37 #define FACETS_HASH_FUNCTION(src, len) XXH3_64bits(src, len)
38 #define FACETS_HASH_ZERO (FACETS_HASH)0
39 -#define FACETS_HASH_UNSAMPLED (FACETS_HASH)UINT64_MAX
39 +#define FACETS_HASH_UNSAMPLED (FACETS_HASH)(UINT64_MAX - 1)
40 +#define FACETS_HASH_ESTIMATED (FACETS_HASH)UINT64_MAX
41
42 static inline void facets_hash_to_str(FACETS_HASH num, char *out) {
43 out[11] = '\0';
@@ -191,11 +192,13 @@ static void simple_hashtable_resize_double(SIMPLE_HASHTABLE *ht) {
192 typedef struct facet_value {
193 FACETS_HASH hash;
194 const char *name;
195 + const char *color;
196 uint32_t name_len;
197
198 bool selected;
199 bool empty;
200 bool unsampled;
201 + bool estimated;
202
203 uint32_t rows_matching_facet_value;
204 uint32_t final_facet_value_counter;
@@ -212,13 +215,15 @@ typedef enum {
215 FACET_KEY_VALUE_UPDATED = (1 << 0),
216 FACET_KEY_VALUE_EMPTY = (1 << 1),
217 FACET_KEY_VALUE_UNSAMPLED = (1 << 2),
215 - FACET_KEY_VALUE_COPIED = (1 << 3),
218 + FACET_KEY_VALUE_ESTIMATED = (1 << 3),
219 + FACET_KEY_VALUE_COPIED = (1 << 4),
220 } FACET_KEY_VALUE_FLAGS;
221
222 #define facet_key_value_updated(k) ((k)->current_value.flags & FACET_KEY_VALUE_UPDATED)
223 #define facet_key_value_empty(k) ((k)->current_value.flags & FACET_KEY_VALUE_EMPTY)
224 #define facet_key_value_unsampled(k) ((k)->current_value.flags & FACET_KEY_VALUE_UNSAMPLED)
221 -#define facet_key_value_empty_or_unsampled(k) ((k)->current_value.flags & (FACET_KEY_VALUE_EMPTY|FACET_KEY_VALUE_UNSAMPLED))
225 +#define facet_key_value_estimated(k) ((k)->current_value.flags & FACET_KEY_VALUE_ESTIMATED)
226 +#define facet_key_value_empty_or_unsampled_or_estimated(k) ((k)->current_value.flags & (FACET_KEY_VALUE_EMPTY|FACET_KEY_VALUE_UNSAMPLED|FACET_KEY_VALUE_ESTIMATED))
227 #define facet_key_value_copied(k) ((k)->current_value.flags & FACET_KEY_VALUE_COPIED)
228
229 struct facet_key {
@@ -260,6 +265,10 @@ struct facet_key {
265 FACET_VALUE *v;
266 } unsampled_value;
267
268 + struct {
269 + FACET_VALUE *v;
270 + } estimated_value;
271 +
272 struct {
273 facet_dynamic_row_t cb;
274 void *data;
@@ -359,6 +368,7 @@ struct facets {
368 size_t evaluated;
369 size_t matched;
370 size_t unsampled;
371 + size_t estimated;
372 size_t created;
373 size_t reused;
374 } rows;
@@ -374,6 +384,7 @@ struct facets {
384 size_t dynamic;
385 size_t empty;
386 size_t unsampled;
387 + size_t estimated;
388 size_t indexed;
389 size_t inserts;
390 size_t conflicts;
@@ -383,6 +394,10 @@ struct facets {
394 size_t searches;
395 } fts;
396 } operations;
397 +
398 + struct {
399 + DICTIONARY *used_hashes_registry;
400 + } report;
401 };
402
403 usec_t facets_row_oldest_ut(FACETS *facets) {
@@ -507,7 +522,17 @@ static inline FACET_VALUE *FACET_VALUE_ADD_TO_INDEX(FACET_KEY *k, const FACET_VA
522
523 memcpy(v, tv, sizeof(*v));
524
510 - DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(k->values.ll, v, prev, next);
525 + if(v->estimated || v->unsampled) {
526 + if(k->values.ll && k->values.ll->estimated) {
527 + FACET_VALUE *estimated = k->values.ll;
528 + DOUBLE_LINKED_LIST_INSERT_ITEM_AFTER_UNSAFE(k->values.ll, estimated, v, prev, next);
529 + }
530 + else
531 + DOUBLE_LINKED_LIST_PREPEND_ITEM_UNSAFE(k->values.ll, v, prev, next);
532 + }
533 + else
534 + DOUBLE_LINKED_LIST_APPEND_ITEM_UNSAFE(k->values.ll, v, prev, next);
535 +
536 k->values.used++;
537
538 if(!v->selected)
@@ -533,6 +558,8 @@ static inline void FACET_VALUE_ADD_UNSAMPLED_VALUE_TO_INDEX(FACET_KEY *k) {
558 .hash = FACETS_HASH_UNSAMPLED,
559 .name = FACET_VALUE_UNSAMPLED,
560 .name_len = sizeof(FACET_VALUE_UNSAMPLED) - 1,
561 + .unsampled = true,
562 + .color = "offline",
563 };
564
565 k->current_value.hash = FACETS_HASH_UNSAMPLED;
@@ -549,11 +576,35 @@ static inline void FACET_VALUE_ADD_UNSAMPLED_VALUE_TO_INDEX(FACET_KEY *k) {
576 }
577 }
578
579 +static inline void FACET_VALUE_ADD_ESTIMATED_VALUE_TO_INDEX(FACET_KEY *k) {
580 + static const FACET_VALUE tv = {
581 + .hash = FACETS_HASH_ESTIMATED,
582 + .name = FACET_VALUE_ESTIMATED,
583 + .name_len = sizeof(FACET_VALUE_ESTIMATED) - 1,
584 + .estimated = true,
585 + .color = "generic",
586 + };
587 +
588 + k->current_value.hash = FACETS_HASH_ESTIMATED;
589 +
590 + if(k->estimated_value.v) {
591 + FACET_VALUE_ADD_CONFLICT(k, k->estimated_value.v, &tv);
592 + k->current_value.v = k->estimated_value.v;
593 + }
594 + else {
595 + FACET_VALUE *v = FACET_VALUE_ADD_TO_INDEX(k, &tv);
596 + v->estimated = true;
597 + k->estimated_value.v = v;
598 + k->current_value.v = v;
599 + }
600 +}
601 +
602 static inline void FACET_VALUE_ADD_EMPTY_VALUE_TO_INDEX(FACET_KEY *k) {
603 static const FACET_VALUE tv = {
604 .hash = FACETS_HASH_ZERO,
605 .name = FACET_VALUE_UNSET,
606 .name_len = sizeof(FACET_VALUE_UNSET) - 1,
607 + .empty = true,
608 };
609
610 k->current_value.hash = FACETS_HASH_ZERO;
@@ -578,6 +629,9 @@ static inline void FACET_VALUE_ADD_CURRENT_VALUE_TO_INDEX(FACET_KEY *k) {
629 tv.name = facets_key_get_value(k);
630 tv.name_len = facets_key_get_value_length(k);
631 tv.hash = FACETS_HASH_FUNCTION(tv.name, tv.name_len);
632 + tv.empty = false;
633 + tv.estimated = false;
634 + tv.unsampled = false;
635
636 k->current_value.v = FACET_VALUE_ADD_TO_INDEX(k, &tv);
637 k->facets->operations.values.indexed++;
@@ -763,6 +817,10 @@ bool facets_key_name_is_facet(FACETS *facets, const char *key) {
817
818 // ----------------------------------------------------------------------------
819
820 +size_t facets_histogram_slots(FACETS *facets) {
821 + return facets->histogram.slots;
822 +}
823 +
824 static usec_t calculate_histogram_bar_width(usec_t after_ut, usec_t before_ut) {
825 // Array of valid durations in seconds
826 static time_t valid_durations_s[] = {
@@ -792,10 +850,6 @@ static inline usec_t facets_histogram_slot_baseline_ut(FACETS *facets, usec_t ut
850 return ut - delta_ut;
851 }
852
795 -size_t facets_histogram_slots(FACETS *facets) {
796 - return facets->histogram.slots;
797 -}
798 -
853 void facets_set_timeframe_and_histogram_by_id(FACETS *facets, const char *key_id, usec_t after_ut, usec_t before_ut) {
854 if(after_ut > before_ut) {
855 usec_t t = after_ut;
@@ -839,7 +893,7 @@ void facets_set_timeframe_and_histogram_by_name(FACETS *facets, const char *key_
893 facets_set_timeframe_and_histogram_by_id(facets, hash_str, after_ut, before_ut);
894 }
895
842 -static inline void facets_histogram_update_value_slot(FACETS *facets, usec_t usec, FACET_VALUE *v) {
896 +static inline uint32_t facets_histogram_slot_at_time_ut(FACETS *facets, usec_t usec, FACET_VALUE *v) {
897 if(unlikely(!v->histogram))
898 v->histogram = callocz(facets->histogram.slots, sizeof(*v->histogram));
899
@@ -856,6 +910,11 @@ static inline void facets_histogram_update_value_slot(FACETS *facets, usec_t use
910 if(unlikely(slot >= facets->histogram.slots))
911 slot = facets->histogram.slots - 1;
912
913 + return slot;
914 +}
915 +
916 +static inline void facets_histogram_update_value_slot(FACETS *facets, usec_t usec, FACET_VALUE *v) {
917 + uint32_t slot = facets_histogram_slot_at_time_ut(facets, usec, v);
918 v->histogram[slot]++;
919 }
920
@@ -872,6 +931,54 @@ static inline void facets_histogram_update_value(FACETS *facets, usec_t usec) {
931 facets_histogram_update_value_slot(facets, usec, v);
932 }
933
934 +void facets_update_estimations(FACETS *facets, usec_t from_ut, usec_t to_ut, size_t entries) {
935 + facets->operations.rows.evaluated += entries;
936 + facets->operations.rows.matched += entries;
937 + facets->operations.rows.estimated += entries;
938 +
939 + if (!facets->histogram.enabled ||
940 + !facets->histogram.key ||
941 + !facets->histogram.key->values.enabled)
942 + return;
943 +
944 + if (from_ut < facets->histogram.after_ut)
945 + from_ut = facets->histogram.after_ut;
946 +
947 + if (to_ut > facets->histogram.before_ut)
948 + to_ut = facets->histogram.before_ut;
949 +
950 + if (!facets->histogram.key->estimated_value.v)
951 + FACET_VALUE_ADD_ESTIMATED_VALUE_TO_INDEX(facets->histogram.key);
952 +
953 + FACET_VALUE *v = facets->histogram.key->estimated_value.v;
954 +
955 + size_t from_slot = facets_histogram_slot_at_time_ut(facets, from_ut, v);
956 + size_t to_slot = facets_histogram_slot_at_time_ut(facets, to_ut, v);
957 + size_t total_ut = to_ut - from_ut;
958 + ssize_t remaining_entries = (ssize_t)entries;
959 +
960 + for (size_t slot = from_slot; slot <= to_slot; slot++) {
961 + if (unlikely(slot >= facets->histogram.slots))
962 + break;
963 +
964 + usec_t slot_start_ut = facets->histogram.after_ut + slot * facets->histogram.slot_width_ut;
965 + usec_t slot_end_ut = slot_start_ut + facets->histogram.slot_width_ut;
966 + usec_t overlap_start_ut = (from_ut > slot_start_ut) ? from_ut : slot_start_ut;
967 + usec_t overlap_end_ut = (to_ut < slot_end_ut) ? to_ut : slot_end_ut;
968 + usec_t overlap_ut = (overlap_end_ut > overlap_start_ut) ? (overlap_end_ut - overlap_start_ut) : 0;
969 +
970 + size_t slot_entries = (overlap_ut * entries) / total_ut;
971 + v->histogram[slot] += slot_entries;
972 + remaining_entries -= (ssize_t)slot_entries;
973 + }
974 +
975 + // Check if all entries are assigned
976 + // This should always be true if the distribution is correct
977 + internal_fatal(remaining_entries < 0 || remaining_entries > (ssize_t)(to_slot - from_slot),
978 + "distribution of estimations is not accurate - there are %zd remaining entries",
979 + remaining_entries);
980 +}
981 +
982 void facets_row_finished_unsampled(FACETS *facets, usec_t usec) {
983 facets->operations.rows.evaluated++;
984 facets->operations.rows.matched++;
@@ -893,8 +1000,72 @@ void facets_row_finished_unsampled(FACETS *facets, usec_t usec) {
1000 facets_reset_key(facets->histogram.key);
1001 }
1002
1003 +static const char *facets_key_name_cached(FACET_KEY *k, DICTIONARY *used_hashes_registry) {
1004 + if(k->name) {
1005 + if(used_hashes_registry && !k->default_selected_for_values) {
1006 + char hash_str[FACET_STRING_HASH_SIZE];
1007 + facets_hash_to_str(k->hash, hash_str);
1008 + dictionary_set(used_hashes_registry, hash_str, (void *)k->name, strlen(k->name) + 1);
1009 + }
1010 +
1011 + return k->name;
1012 + }
1013 +
1014 + // key has no name
1015 + const char *name = "[UNAVAILABLE_FIELD]";
1016 +
1017 + if(used_hashes_registry) {
1018 + char hash_str[FACET_STRING_HASH_SIZE];
1019 + facets_hash_to_str(k->hash, hash_str);
1020 + const char *s = dictionary_get(used_hashes_registry, hash_str);
1021 + if(s) name = s;
1022 + }
1023 +
1024 + return name;
1025 +}
1026 +
1027 +static const char *facets_key_value_cached(FACET_KEY *k, FACET_VALUE *v, DICTIONARY *used_hashes_registry) {
1028 + if(v->empty || v->estimated || v->unsampled)
1029 + return v->name;
1030 +
1031 + if(v->name && v->name_len) {
1032 + if(used_hashes_registry && !k->default_selected_for_values && v->selected) {
1033 + char hash_str[FACET_STRING_HASH_SIZE];
1034 + facets_hash_to_str(v->hash, hash_str);
1035 + dictionary_set(used_hashes_registry, hash_str, (void *)v->name, v->name_len + 1);
1036 + }
1037 +
1038 + return v->name;
1039 + }
1040 +
1041 + // key has no name
1042 + const char *name = "[unavailable field]";
1043 +
1044 + if(used_hashes_registry) {
1045 + char hash_str[FACET_STRING_HASH_SIZE];
1046 + facets_hash_to_str(v->hash, hash_str);
1047 + const char *s = dictionary_get(used_hashes_registry, hash_str);
1048 + if(s) name = s;
1049 + }
1050 +
1051 + return name;
1052 +}
1053 +
1054 +static inline void facets_key_value_transformed(FACETS *facets, FACET_KEY *k, FACET_VALUE *v, BUFFER *dst, FACETS_TRANSFORMATION_SCOPE scope) {
1055 + buffer_flush(dst);
1056 +
1057 + if(v->empty || v->unsampled || v->estimated)
1058 + buffer_strcat(dst, v->name);
1059 + else if(k->transform.cb && k->transform.view_only) {
1060 + buffer_contents_replace(dst, v->name, v->name_len);
1061 + k->transform.cb(facets, dst, scope, k->transform.data);
1062 + }
1063 + else
1064 + buffer_strcat(dst, facets_key_value_cached(k, v, facets->report.used_hashes_registry));
1065 +}
1066 +
1067 static inline void facets_histogram_value_names(BUFFER *wb, FACETS *facets __maybe_unused, FACET_KEY *k, const char *key, const char *first_key) {
897 - BUFFER *tb = NULL;
1068 + CLEAN_BUFFER *tb = buffer_create(0, NULL);
1069
1070 buffer_json_member_add_array(wb, key);
1071 {
@@ -907,23 +1078,30 @@ static inline void facets_histogram_value_names(BUFFER *wb, FACETS *facets __may
1078 if (unlikely(!v->histogram))
1079 continue;
1080
910 - if(!v->empty && k->transform.cb && k->transform.view_only) {
911 - if(!tb)
912 - tb = buffer_create(0, NULL);
913 -
914 - buffer_contents_replace(tb, v->name, v->name_len);
915 - k->transform.cb(facets, tb, FACETS_TRANSFORM_HISTOGRAM, k->transform.data);
916 - buffer_json_add_array_item_string(wb, buffer_tostring(tb));
917 - }
918 - else
919 - buffer_json_add_array_item_string(wb, v->name);
1081 + facets_key_value_transformed(facets, k, v, tb, FACETS_TRANSFORM_HISTOGRAM);
1082 + buffer_json_add_array_item_string(wb, buffer_tostring(tb));
1083 }
1084 foreach_value_in_key_done(v);
1085 }
1086 }
1087 buffer_json_array_close(wb); // key
1088 +}
1089
926 - buffer_free(tb);
1090 +static inline void facets_histogram_value_colors(BUFFER *wb, FACETS *facets __maybe_unused, FACET_KEY *k, const char *key) {
1091 + buffer_json_member_add_array(wb, key);
1092 + {
1093 + if(k && k->values.enabled) {
1094 + FACET_VALUE *v;
1095 + foreach_value_in_key(k, v) {
1096 + if (unlikely(!v->histogram))
1097 + continue;
1098 +
1099 + buffer_json_add_array_item_string(wb, v->color);
1100 + }
1101 + foreach_value_in_key_done(v);
1102 + }
1103 + }
1104 + buffer_json_array_close(wb); // key
1105 }
1106
1107 static inline void facets_histogram_value_units(BUFFER *wb, FACETS *facets __maybe_unused, FACET_KEY *k, const char *key) {
@@ -1029,6 +1207,8 @@ static inline void facets_histogram_value_con(BUFFER *wb, FACETS *facets __maybe
1207 }
1208
1209 static void facets_histogram_generate(FACETS *facets, FACET_KEY *k, BUFFER *wb) {
1210 + CLEAN_BUFFER *tmp = buffer_create(0, NULL);
1211 +
1212 size_t dimensions = 0;
1213 uint32_t min = UINT32_MAX, max = 0, sum = 0, count = 0;
1214
@@ -1070,6 +1250,7 @@ static void facets_histogram_generate(FACETS *facets, FACET_KEY *k, BUFFER *wb)
1250
1251 buffer_json_member_add_object(wb, "summary");
1252 {
1253 + // summary.nodes
1254 buffer_json_member_add_array(wb, "nodes");
1255 {
1256 buffer_json_add_array_item_object(wb); // node
@@ -1116,6 +1297,7 @@ static void facets_histogram_generate(FACETS *facets, FACET_KEY *k, BUFFER *wb)
1297 }
1298 buffer_json_array_close(wb); // nodes
1299
1300 + // summary.contexts
1301 buffer_json_member_add_array(wb, "contexts");
1302 {
1303 buffer_json_add_array_item_object(wb); // context
@@ -1153,6 +1335,7 @@ static void facets_histogram_generate(FACETS *facets, FACET_KEY *k, BUFFER *wb)
1335 }
1336 buffer_json_array_close(wb); // contexts
1337
1338 + // summary.instances
1339 buffer_json_member_add_array(wb, "instances");
1340 {
1341 buffer_json_add_array_item_object(wb); // instance
@@ -1184,17 +1367,20 @@ static void facets_histogram_generate(FACETS *facets, FACET_KEY *k, BUFFER *wb)
1367 }
1368 buffer_json_array_close(wb); // instances
1369
1370 + // summary.dimensions
1371 buffer_json_member_add_array(wb, "dimensions");
1372 if(dimensions && k && k->values.enabled) {
1373 size_t pri = 0;
1374 FACET_VALUE *v;
1375 +
1376 foreach_value_in_key(k, v) {
1377 if(unlikely(!v->histogram))
1378 continue;
1379
1380 buffer_json_add_array_item_object(wb); // dimension
1381 {
1197 - buffer_json_member_add_string(wb, "id", v->name);
1382 + facets_key_value_transformed(facets, k, v, tmp, FACETS_TRANSFORM_HISTOGRAM);
1383 + buffer_json_member_add_string(wb, "id", buffer_tostring(tmp));
1384 buffer_json_member_add_object(wb, "ds");
1385 {
1386 buffer_json_member_add_uint64(wb, "sl", 1);
@@ -1368,6 +1554,7 @@ static void facets_histogram_generate(FACETS *facets, FACET_KEY *k, BUFFER *wb)
1554
1555 facets_histogram_value_names(wb, facets, k, "ids", NULL);
1556 facets_histogram_value_names(wb, facets, k, "names", NULL);
1557 + facets_histogram_value_colors(wb, facets, k, "colors");
1558 facets_histogram_value_units(wb, facets, k, "units");
1559
1560 buffer_json_member_add_object(wb, "sts");
@@ -1653,7 +1840,7 @@ static inline void facets_key_check_value(FACETS *facets, FACET_KEY *k) {
1840 facets->keys_in_row.array[facets->keys_in_row.used++] = k;
1841
1842 k->current_value.flags |= FACET_KEY_VALUE_UPDATED;
1656 - k->current_value.flags &= ~(FACET_KEY_VALUE_EMPTY|FACET_KEY_VALUE_UNSAMPLED);
1843 + k->current_value.flags &= ~(FACET_KEY_VALUE_EMPTY|FACET_KEY_VALUE_UNSAMPLED|FACET_KEY_VALUE_ESTIMATED);
1844
1845 facets->operations.values.registered++;
1846
@@ -1667,7 +1854,7 @@ static inline void facets_key_check_value(FACETS *facets, FACET_KEY *k) {
1854 // if(strstr(buffer_tostring(k->current_value), "fprintd") != NULL)
1855 // found = true;
1856
1670 - if(facets->query && !facet_key_value_empty_or_unsampled(k) && ((k->options & FACET_KEY_OPTION_FTS) || facets->options & FACETS_OPTION_ALL_KEYS_FTS)) {
1857 + if(facets->query && !facet_key_value_empty_or_unsampled_or_estimated(k) && ((k->options & FACET_KEY_OPTION_FTS) || facets->options & FACETS_OPTION_ALL_KEYS_FTS)) {
1858 facets->operations.fts.searches++;
1859 facets_key_value_copy_to_buffer(k);
1860 switch(simple_pattern_matches_extract(facets->query, buffer_tostring(k->current_value.b), NULL, 0)) {
@@ -1778,7 +1965,7 @@ static FACET_ROW *facets_row_create(FACETS *facets, usec_t usec, FACET_ROW *into
1965 .empty = true,
1966 };
1967
1781 - if(facet_key_value_updated(k) && !facet_key_value_empty_or_unsampled(k)) {
1968 + if(facet_key_value_updated(k) && !facet_key_value_empty_or_unsampled_or_estimated(k)) {
1969 t.tmp = facets_key_get_value(k);
1970 t.tmp_len = facets_key_get_value_length(k);
1971 t.empty = false;
@@ -2210,7 +2397,7 @@ static uint32_t facets_sort_and_reorder_values(FACET_KEY *k) {
2397 if(!k->values.enabled || !k->values.ll || !k->values.used)
2398 return 0;
2399
2213 - if(!k->transform.cb || !(k->facets->options & FACETS_OPTION_SORT_FACETS_ALPHABETICALLY))
2400 + if(!k->transform.cb || !k->transform.view_only || !(k->facets->options & FACETS_OPTION_SORT_FACETS_ALPHABETICALLY))
2401 return facets_sort_and_reorder_values_internal(k);
2402
2403 // we have a transformation and has to be sorted alphabetically
@@ -2234,8 +2421,7 @@ static uint32_t facets_sort_and_reorder_values(FACET_KEY *k) {
2421 values[used].name_len = v->name_len;
2422 used++;
2423
2237 - buffer_contents_replace(tb, v->name, v->name_len);
2238 - k->transform.cb(k->facets, tb, FACETS_TRANSFORM_FACET_SORT, k->transform.data);
2424 + facets_key_value_transformed(k->facets, k, v, tb, FACETS_TRANSFORM_FACET_SORT);
2425 v->name = strdupz(buffer_tostring(tb));
2426 v->name_len = buffer_strlen(tb);
2427 }
@@ -2273,55 +2459,9 @@ void facets_table_config(BUFFER *wb) {
2459 buffer_json_object_close(wb); // pagination
2460 }
2461
2276 -static const char *facets_json_key_name_string(FACET_KEY *k, DICTIONARY *used_hashes_registry) {
2277 - if(k->name) {
2278 - if(used_hashes_registry && !k->default_selected_for_values) {
2279 - char hash_str[FACET_STRING_HASH_SIZE];
2280 - facets_hash_to_str(k->hash, hash_str);
2281 - dictionary_set(used_hashes_registry, hash_str, (void *)k->name, strlen(k->name) + 1);
2282 - }
2283 -
2284 - return k->name;
2285 - }
2286 -
2287 - // key has no name
2288 - const char *name = "[UNAVAILABLE_FIELD]";
2289 -
2290 - if(used_hashes_registry) {
2291 - char hash_str[FACET_STRING_HASH_SIZE];
2292 - facets_hash_to_str(k->hash, hash_str);
2293 - const char *s = dictionary_get(used_hashes_registry, hash_str);
2294 - if(s) name = s;
2295 - }
2296 -
2297 - return name;
2298 -}
2299 -
2300 -static const char *facets_json_key_value_string(FACET_KEY *k, FACET_VALUE *v, DICTIONARY *used_hashes_registry) {
2301 - if(v->name && v->name_len) {
2302 - if(used_hashes_registry && !k->default_selected_for_values && v->selected) {
2303 - char hash_str[FACET_STRING_HASH_SIZE];
2304 - facets_hash_to_str(v->hash, hash_str);
2305 - dictionary_set(used_hashes_registry, hash_str, (void *)v->name, v->name_len + 1);
2306 - }
2307 -
2308 - return v->name;
2309 - }
2310 -
2311 - // key has no name
2312 - const char *name = "[unavailable field]";
2313 -
2314 - if(used_hashes_registry) {
2315 - char hash_str[FACET_STRING_HASH_SIZE];
2316 - facets_hash_to_str(v->hash, hash_str);
2317 - const char *s = dictionary_get(used_hashes_registry, hash_str);
2318 - if(s) name = s;
2319 - }
2320 -
2321 - return name;
2322 -}
2323 -
2462 void facets_report(FACETS *facets, BUFFER *wb, DICTIONARY *used_hashes_registry) {
2463 + facets->report.used_hashes_registry = used_hashes_registry;
2464 +
2465 if(!(facets->options & FACETS_OPTION_DATA_ONLY)) {
2466 facets_table_config(wb);
2467 facets_accepted_parameters_to_json_array(facets, wb, true);
@@ -2345,7 +2485,7 @@ void facets_report(FACETS *facets, BUFFER *wb, DICTIONARY *used_hashes_registry)
2485 }
2486
2487 if(show_facets) {
2348 - BUFFER *tb = NULL;
2488 + CLEAN_BUFFER *tb = buffer_create(0, NULL);
2489 FACET_KEY *k;
2490 foreach_key_in_facets(facets, k) {
2491 if(!k->values.enabled)
@@ -2358,7 +2498,9 @@ void facets_report(FACETS *facets, BUFFER *wb, DICTIONARY *used_hashes_registry)
2498 buffer_json_add_array_item_object(wb); // key
2499 {
2500 buffer_json_member_add_string(wb, "id", hash_to_static_string(k->hash));
2361 - buffer_json_member_add_string(wb, "name", facets_json_key_name_string(k, used_hashes_registry));
2501 + buffer_json_member_add_string(wb, "name", facets_key_name_cached(k
2502 + , facets->report.used_hashes_registry
2503 + ));
2504
2505 if(!k->order) k->order = facets->order++;
2506 buffer_json_member_add_uint64(wb, "order", k->order);
@@ -2370,21 +2512,15 @@ void facets_report(FACETS *facets, BUFFER *wb, DICTIONARY *used_hashes_registry)
2512 if((facets->options & FACETS_OPTION_DONT_SEND_EMPTY_VALUE_FACETS) && v->empty)
2513 continue;
2514
2515 + if(v->unsampled || v->estimated)
2516 + continue;
2517 +
2518 buffer_json_add_array_item_object(wb);
2519 {
2520 buffer_json_member_add_string(wb, "id", hash_to_static_string(v->hash));
2521
2377 - if(!v->empty && k->transform.cb && k->transform.view_only) {
2378 - if(!tb)
2379 - tb = buffer_create(0, NULL);
2380 -
2381 - buffer_contents_replace(tb, v->name, v->name_len);
2382 - k->transform.cb(facets, tb, FACETS_TRANSFORM_FACET, k->transform.data);
2383 - buffer_json_member_add_string(wb, "name", buffer_tostring(tb));
2384 - }
2385 - else
2386 - buffer_json_member_add_string(wb, "name", facets_json_key_value_string(k, v, used_hashes_registry));
2387 -
2522 + facets_key_value_transformed(facets, k, v, tb, FACETS_TRANSFORM_FACET);
2523 + buffer_json_member_add_string(wb, "name", buffer_tostring(tb));
2524 buffer_json_member_add_uint64(wb, "count", v->final_facet_value_counter);
2525 buffer_json_member_add_uint64(wb, "order", v->order);
2526 }
@@ -2397,7 +2533,6 @@ void facets_report(FACETS *facets, BUFFER *wb, DICTIONARY *used_hashes_registry)
2533 buffer_json_object_close(wb); // key
2534 }
2535 foreach_key_in_facets_done(k);
2400 - buffer_free(tb);
2536 buffer_json_array_close(wb); // facets
2537 }
2538 }
@@ -2612,6 +2747,7 @@ void facets_report(FACETS *facets, BUFFER *wb, DICTIONARY *used_hashes_registry)
2747 buffer_json_member_add_uint64(wb, "evaluated", facets->operations.rows.evaluated);
2748 buffer_json_member_add_uint64(wb, "matched", facets->operations.rows.matched);
2749 buffer_json_member_add_uint64(wb, "unsampled", facets->operations.rows.unsampled);
2750 + buffer_json_member_add_uint64(wb, "estimated", facets->operations.rows.estimated);
2751 buffer_json_member_add_uint64(wb, "returned", facets->items_to_return);
2752 buffer_json_member_add_uint64(wb, "max_to_return", facets->max_items_to_return);
2753 buffer_json_member_add_uint64(wb, "before", facets->operations.skips_before);
@@ -2676,6 +2812,8 @@ void facets_report(FACETS *facets, BUFFER *wb, DICTIONARY *used_hashes_registry)
2812 buffer_json_member_add_uint64(wb, "transformed", facets->operations.values.transformed);
2813 buffer_json_member_add_uint64(wb, "dynamic", facets->operations.values.dynamic);
2814 buffer_json_member_add_uint64(wb, "empty", facets->operations.values.empty);
2815 + buffer_json_member_add_uint64(wb, "unsampled", facets->operations.values.unsampled);
2816 + buffer_json_member_add_uint64(wb, "estimated", facets->operations.values.estimated);
2817 buffer_json_member_add_uint64(wb, "indexed", facets->operations.values.indexed);
2818 buffer_json_member_add_uint64(wb, "inserts", facets->operations.values.inserts);
2819 buffer_json_member_add_uint64(wb, "conflicts", facets->operations.values.conflicts);
libnetdata/facets/facets.h
+2
@@ -7,6 +7,7 @@
7
8 #define FACET_VALUE_UNSET "-"
9 #define FACET_VALUE_UNSAMPLED "[unsampled]"
10 +#define FACET_VALUE_ESTIMATED "[estimated]"
11
12 typedef enum __attribute__((packed)) {
13 FACETS_ANCHOR_DIRECTION_FORWARD,
@@ -87,6 +88,7 @@ void facets_rows_begin(FACETS *facets);
88 bool facets_row_finished(FACETS *facets, usec_t usec);
89
90 void facets_row_finished_unsampled(FACETS *facets, usec_t usec);
91 +void facets_update_estimations(FACETS *facets, usec_t from_ut, usec_t to_ut, size_t entries);
92 size_t facets_histogram_slots(FACETS *facets);
93
94 FACET_KEY *facets_register_key_name(FACETS *facets, const char *key, FACET_KEY_OPTIONS options);
libnetdata/log/log.c
+1 -1
@@ -221,7 +221,7 @@ int nd_log_priority2id(const char *priority) {
221 return NDLP_INFO;
222 }
223
224 -static const char *nd_log_id2priority(ND_LOG_FIELD_PRIORITY priority) {
224 +const char *nd_log_id2priority(ND_LOG_FIELD_PRIORITY priority) {
225 size_t entries = sizeof(nd_log_priorities) / sizeof(nd_log_priorities[0]);
226 for(size_t i = 0; i < entries ;i++) {
227 if(priority == nd_log_priorities[i].priority)
libnetdata/log/log.h
+1
@@ -146,6 +146,7 @@ void nd_log_set_thread_source(ND_LOG_SOURCES source);
146 bool nd_log_journal_socket_available(void);
147 ND_LOG_FIELD_ID nd_log_field_id_by_name(const char *field, size_t len);
148 int nd_log_priority2id(const char *priority);
149 +const char *nd_log_id2priority(ND_LOG_FIELD_PRIORITY priority);
150
151 typedef bool (*log_formatter_callback_t)(BUFFER *wb, void *data);
152