journal: fix incremental queries (#16098)
fix if_modified_since to not return empty responses and no overlapping data
Costa Tsaousis committed
Oct 3, 2023 at 11:57 UTC
20e3113847f8d525dc9c1b034e73ce1fb753207b
2 files changed
+45
-35
collectors/systemd-journal.plugin/systemd-journal.c
+24
-14
@@ -330,7 +330,8 @@ static inline size_t netdata_systemd_journal_process_row(sd_journal *j, FACETS *
330
331
#define FUNCTION_PROGRESS_UPDATE_ROWS(rows_read, rows) __atomic_fetch_add(&(rows_read), rows, __ATOMIC_RELAXED)
332
#define FUNCTION_PROGRESS_UPDATE_BYTES(bytes_read, bytes) __atomic_fetch_add(&(bytes_read), bytes, __ATOMIC_RELAXED)
333
-#define FUNCTION_PROGRESS_EVERY_ROWS 10000
333
+#define FUNCTION_PROGRESS_EVERY_ROWS (1ULL << 13)
334
+#define FUNCTION_DATA_ONLY_CHECK_EVERY_ROWS (1ULL << 7)
335
336
static inline ND_SD_JOURNAL_STATUS check_stop(const bool *cancelled, const usec_t *stop_monotonic_ut) {
337
if(cancelled && __atomic_load_n(cancelled, __ATOMIC_RELAXED)) {
@@ -352,7 +353,7 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_backward(
353
354
usec_t anchor_delta = __atomic_load_n(&jf->max_journal_vs_realtime_delta_ut, __ATOMIC_RELAXED);
355
355
- usec_t start_ut = ((fqs->data_only && fqs->anchor.start_ut) ? fqs->anchor.start_ut : fqs->before_ut) + anchor_delta;
356
+ usec_t start_ut = ((fqs->data_only && fqs->anchor.start_ut) ? fqs->anchor.start_ut : fqs->before_ut + USEC_PER_SEC - 1) + anchor_delta;
357
usec_t stop_ut = (fqs->data_only && fqs->anchor.stop_ut) ? fqs->anchor.stop_ut : fqs->after_ut;
358
359
if(!netdata_systemd_journal_seek_to(j, start_ut))
@@ -360,7 +361,7 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_backward(
361
362
size_t errors_no_timestamp = 0;
363
usec_t earliest_msg_ut = 0;
363
- size_t row_counter = 0, last_row_counter = 0;
364
+ size_t row_counter = 0, last_row_counter = 0, rows_useful = 0;
365
size_t bytes = 0, last_bytes = 0;
366
367
ND_SD_JOURNAL_STATUS status = ND_SD_JOURNAL_OK;
@@ -384,10 +385,10 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_backward(
385
386
bytes += netdata_systemd_journal_process_row(j, facets, jf, &msg_ut);
387
if(facets_row_finished(facets, msg_ut))
387
- fqs->rows_useful++;
388
+ rows_useful++;
389
390
row_counter++;
390
- if(row_counter % 100 == 0 && fqs->data_only && facets_rows(facets) >= fqs->entries) {
391
+ if(row_counter % FUNCTION_DATA_ONLY_CHECK_EVERY_ROWS == 0 && fqs->data_only && facets_rows(facets) >= fqs->entries) {
392
// stop the data only query
393
usec_t oldest = facets_row_oldest_ut(facets);
394
if(oldest && msg_ut < (oldest - anchor_delta))
@@ -408,6 +409,8 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_backward(
409
FUNCTION_PROGRESS_UPDATE_ROWS(fqs->rows_read, row_counter - last_row_counter);
410
FUNCTION_PROGRESS_UPDATE_BYTES(fqs->bytes_read, bytes - last_bytes);
411
412
+ fqs->rows_useful += rows_useful;
413
+
414
if(errors_no_timestamp)
415
netdata_log_error("SYSTEMD-JOURNAL: %zu lines did not have timestamps", errors_no_timestamp);
416
@@ -424,14 +427,14 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_forward(
427
usec_t anchor_delta = __atomic_load_n(&jf->max_journal_vs_realtime_delta_ut, __ATOMIC_RELAXED);
428
429
usec_t start_ut = (fqs->data_only && fqs->anchor.start_ut) ? fqs->anchor.start_ut : fqs->after_ut;
427
- usec_t stop_ut = ((fqs->data_only && fqs->anchor.stop_ut) ? fqs->anchor.stop_ut : fqs->before_ut) + anchor_delta;
430
+ usec_t stop_ut = ((fqs->data_only && fqs->anchor.stop_ut) ? fqs->anchor.stop_ut : fqs->before_ut + USEC_PER_SEC - 1) + anchor_delta;
431
432
if(!netdata_systemd_journal_seek_to(j, start_ut))
433
return ND_SD_JOURNAL_FAILED_TO_SEEK;
434
435
size_t errors_no_timestamp = 0;
436
usec_t earliest_msg_ut = 0;
434
- size_t row_counter = 0, last_row_counter = 0;
437
+ size_t row_counter = 0, last_row_counter = 0, rows_useful = 0;
438
size_t bytes = 0, last_bytes = 0;
439
440
ND_SD_JOURNAL_STATUS status = ND_SD_JOURNAL_OK;
@@ -455,10 +458,10 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_forward(
458
459
bytes += netdata_systemd_journal_process_row(j, facets, jf, &msg_ut);
460
if(facets_row_finished(facets, msg_ut))
458
- fqs->rows_useful++;
461
+ rows_useful++;
462
463
row_counter++;
461
- if(row_counter % 100 == 0 && fqs->data_only && facets_rows(facets) >= fqs->entries) {
464
+ if(row_counter % FUNCTION_DATA_ONLY_CHECK_EVERY_ROWS == 0 && fqs->data_only && facets_rows(facets) >= fqs->entries) {
465
usec_t newest = facets_row_newest_ut(facets);
466
if(newest && msg_ut > (newest + anchor_delta))
467
break;
@@ -478,6 +481,8 @@ ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_forward(
481
FUNCTION_PROGRESS_UPDATE_ROWS(fqs->rows_read, row_counter - last_row_counter);
482
FUNCTION_PROGRESS_UPDATE_BYTES(fqs->bytes_read, bytes - last_bytes);
483
484
+ fqs->rows_useful += rows_useful;
485
+
486
if(errors_no_timestamp)
487
netdata_log_error("SYSTEMD-JOURNAL: %zu lines did not have timestamps", errors_no_timestamp);
488
@@ -1093,6 +1098,9 @@ static int netdata_systemd_journal_query(BUFFER *wb, FACETS *facets, FUNCTION_QU
1098
1099
fqs->files_matched = 0;
1100
fqs->file_working = 0;
1101
+ fqs->rows_useful = 0;
1102
+ fqs->rows_read = 0;
1103
+ fqs->bytes_read = 0;
1104
1105
size_t files_used = 0;
1106
size_t files_max = dictionary_entries(journal_files_registry);
@@ -1118,10 +1126,6 @@ static int netdata_systemd_journal_query(BUFFER *wb, FACETS *facets, FUNCTION_QU
1126
return HTTP_RESP_NOT_MODIFIED;
1127
}
1128
1121
- // We will not do an if_modified_since query
1122
- // we know something changed in the files
1123
- fqs->if_modified_since = 0;
1124
-
1129
// sort the files, so that they are optimal for facets
1130
if(files_used >= 2) {
1131
if (fqs->direction == FACETS_ANCHOR_DIRECTION_BACKWARD)
@@ -1222,6 +1226,12 @@ static int netdata_systemd_journal_query(BUFFER *wb, FACETS *facets, FUNCTION_QU
1226
1227
switch (status) {
1228
case ND_SD_JOURNAL_OK:
1229
+ if(fqs->if_modified_since && !fqs->rows_useful) {
1230
+ buffer_flush(wb);
1231
+ return HTTP_RESP_NOT_MODIFIED;
1232
+ }
1233
+ break;
1234
+
1235
case ND_SD_JOURNAL_TIMED_OUT:
1236
case ND_SD_JOURNAL_NO_FILE_MATCHED:
1237
break;
@@ -2575,7 +2585,7 @@ int main(int argc __maybe_unused, char **argv __maybe_unused) {
2585
2586
if(argc == 2 && strcmp(argv[1], "debug") == 0) {
2587
bool cancelled = false;
2578
- char buf[] = "systemd-journal after:-2592000 before:0 last:500";
2588
+ char buf[] = "systemd-journal after:1696319393 before:1696320293 anchor:1696320283039944 direction:forward last:100 if_modified_since:1696320283039989 data_only:true delta:true tail:true slice:true source:all histogram:DHKucpqUoe1";
2589
// char buf[] = "systemd-journal after:1695332964 before:1695937764 direction:backward last:100 slice:true source:all DHKucpqUoe1:PtVoyIuX.MU";
2590
// char buf[] = "systemd-journal after:1694511062 before:1694514662 anchor:1694514122024403";
2591
function_systemd_journal("123", buf, 600, &cancelled);
libnetdata/facets/facets.c
+21
-21
@@ -1772,40 +1772,40 @@ bool facets_row_finished(FACETS *facets, usec_t usec) {
1772
1773
bool within_anchor = facets_is_entry_within_anchor(facets, usec);
1774
1775
- if(within_anchor && selected_keys >= total_keys - 1) {
1776
- size_t found = 0; (void)found;
1775
+ if(likely(within_anchor)) {
1776
+ if(selected_keys >= total_keys - 1) {
1777
+ size_t found = 0;
1778
+ (void) found;
1779
1778
- for(size_t p = 0; p < entries ;p++) {
1779
- FACET_KEY *k = facets->keys_with_values.array[p];
1780
+ for(size_t p = 0; p < entries; p++) {
1781
+ FACET_KEY *k = facets->keys_with_values.array[p];
1782
1781
- size_t counted_by = selected_keys;
1783
+ size_t counted_by = selected_keys;
1784
1783
- if (counted_by != total_keys && !k->key_values_selected_in_row)
1784
- counted_by++;
1785
+ if(counted_by != total_keys && !k->key_values_selected_in_row)
1786
+ counted_by++;
1787
1786
- if (counted_by == total_keys) {
1787
- FACET_VALUE *v = FACET_VALUE_GET_CURRENT_VALUE(k);
1788
- v->final_facet_value_counter++;
1788
+ if(counted_by == total_keys) {
1789
+ FACET_VALUE *v = FACET_VALUE_GET_CURRENT_VALUE(k);
1790
+ v->final_facet_value_counter++;
1791
1790
- found++;
1792
+ found++;
1793
+ }
1794
}
1792
- }
1793
-
1794
- internal_fatal(!found, "We should find at least one facet to count this row");
1795
- }
1795
1797
- if(selected_keys == total_keys) {
1798
- // we need to keep this row
1799
-
1800
- facets_histogram_update_value(facets, usec);
1796
+ internal_fatal(!found, "We should find at least one facet to count this row");
1797
+ }
1798
1802
- if(within_anchor)
1799
+ if(selected_keys == total_keys) {
1800
+ // we need to keep this row
1801
+ facets_histogram_update_value(facets, usec);
1802
facets_row_keep(facets, usec);
1803
+ }
1804
}
1805
1806
facets_reset_keys_with_value_and_row(facets);
1807
1808
- return selected_keys == total_keys;
1808
+ return selected_keys == total_keys && within_anchor;
1809
}
1810
1811
// ----------------------------------------------------------------------------