@cryptotaxi247 / netdata-1 / commits / 49c6ca778

journal: go up to stop anchor on data only queries (#16107)

go up to stop anchor on data only queries

Costa Tsaousis committed Oct 3, 2023 at 14:38 UTC 49c6ca7787289815833ada2ab5bb93762072e683
1 file changed +15 -8
collectors/systemd-journal.plugin/systemd-journal.c
+15 -8
@@ -97,10 +97,13 @@ int fstat64(int fd, struct stat64 *buf) {
97 #define SYSTEMD_JOURNAL_FUNCTION_NAME "systemd-journal"
98 #define SYSTEMD_JOURNAL_DEFAULT_TIMEOUT 60
99 #define SYSTEMD_JOURNAL_MAX_PARAMS 100
100 -#define SYSTEMD_JOURNAL_DEFAULT_QUERY_DURATION (3 * 3600)
100 +#define SYSTEMD_JOURNAL_DEFAULT_QUERY_DURATION (1 * 3600)
101 #define SYSTEMD_JOURNAL_DEFAULT_ITEMS_PER_QUERY 200
102 #define SYSTEMD_JOURNAL_WORKER_THREADS 5
103
104 +#define JOURNAL_VS_REALTIME_DELTA_DEFAULT_UT (5 * USEC_PER_SEC) // assume always 5 seconds latency
105 +#define JOURNAL_VS_REALTIME_DELTA_MAX_UT (2 * 60 * USEC_PER_SEC) // up to 2 minutes latency
106 +
107 #define JOURNAL_PARAMETER_HELP "help"
108 #define JOURNAL_PARAMETER_AFTER "after"
109 #define JOURNAL_PARAMETER_BEFORE "before"
@@ -278,9 +281,6 @@ static inline bool netdata_systemd_journal_seek_to(sd_journal *j, usec_t timesta
281
282 #define JD_SOURCE_REALTIME_TIMESTAMP "_SOURCE_REALTIME_TIMESTAMP"
283
281 -#define JOURNAL_VS_REALTIME_DELTA_DEFAULT_UT (2 * USEC_PER_SEC) // assume always 2 seconds latency
282 -#define JOURNAL_VS_REALTIME_DELTA_MAX_UT (2 * 60 * USEC_PER_SEC) // up to 2 minutes delta
283 -
284 static inline bool parse_journal_field(const char *data, size_t data_length, const char **key, size_t *key_length, const char **value, size_t *value_length) {
285 const char *k = data;
286 const char *equal = strchr(k, '=');
@@ -371,6 +371,7 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_backward(
371
372 usec_t start_ut = ((fqs->data_only && fqs->anchor.start_ut) ? fqs->anchor.start_ut : fqs->before_ut) + anchor_delta;
373 usec_t stop_ut = (fqs->data_only && fqs->anchor.stop_ut) ? fqs->anchor.stop_ut : fqs->after_ut;
374 + bool stop_when_full = (fqs->data_only && !fqs->anchor.stop_ut);
375
376 if(!netdata_systemd_journal_seek_to(j, start_ut))
377 return ND_SD_JOURNAL_FAILED_TO_SEEK;
@@ -404,14 +405,16 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_backward(
405 rows_useful++;
406
407 row_counter++;
407 - if(row_counter % FUNCTION_DATA_ONLY_CHECK_EVERY_ROWS == 0 && fqs->data_only && facets_rows(facets) >= fqs->entries) {
408 + if(unlikely((row_counter % FUNCTION_DATA_ONLY_CHECK_EVERY_ROWS) == 0 &&
409 + stop_when_full &&
410 + facets_rows(facets) >= fqs->entries)) {
411 // stop the data only query
412 usec_t oldest = facets_row_oldest_ut(facets);
413 if(oldest && msg_ut < (oldest - anchor_delta))
414 break;
415 }
416
414 - if(row_counter % FUNCTION_PROGRESS_EVERY_ROWS == 0) {
417 + if(unlikely(row_counter % FUNCTION_PROGRESS_EVERY_ROWS == 0)) {
418 FUNCTION_PROGRESS_UPDATE_ROWS(fqs->rows_read, row_counter - last_row_counter);
419 last_row_counter = row_counter;
420
@@ -444,6 +447,7 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_forward(
447
448 usec_t start_ut = (fqs->data_only && fqs->anchor.start_ut) ? fqs->anchor.start_ut : fqs->after_ut;
449 usec_t stop_ut = ((fqs->data_only && fqs->anchor.stop_ut) ? fqs->anchor.stop_ut : fqs->before_ut) + anchor_delta;
450 + bool stop_when_full = (fqs->data_only && !fqs->anchor.stop_ut);
451
452 if(!netdata_systemd_journal_seek_to(j, start_ut))
453 return ND_SD_JOURNAL_FAILED_TO_SEEK;
@@ -477,13 +481,16 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_forward(
481 rows_useful++;
482
483 row_counter++;
480 - if(row_counter % FUNCTION_DATA_ONLY_CHECK_EVERY_ROWS == 0 && fqs->data_only && facets_rows(facets) >= fqs->entries) {
484 + if(unlikely((row_counter % FUNCTION_DATA_ONLY_CHECK_EVERY_ROWS) == 0 &&
485 + stop_when_full &&
486 + facets_rows(facets) >= fqs->entries)) {
487 + // stop the data only query
488 usec_t newest = facets_row_newest_ut(facets);
489 if(newest && msg_ut > (newest + anchor_delta))
490 break;
491 }
492
486 - if(row_counter % FUNCTION_PROGRESS_EVERY_ROWS == 0) {
493 + if(unlikely(row_counter % FUNCTION_PROGRESS_EVERY_ROWS == 0)) {
494 FUNCTION_PROGRESS_UPDATE_ROWS(fqs->rows_read, row_counter - last_row_counter);
495 last_row_counter = row_counter;
496