5
* GPL v3+
6
*/
7
8
-#include "systemd-internals.h"
9
-
8
/*
9
* TODO
12
- *
10
* _UDEV_DEVLINK is frequently set more than once per field - support multi-value faces
11
*
12
*/
13
17
-#define FACET_MAX_VALUE_LENGTH 8192
14
+#include "systemd-internals.h"
15
16
#define SYSTEMD_JOURNAL_FUNCTION_DESCRIPTION "View, search and analyze systemd journal entries."
17
#define SYSTEMD_JOURNAL_FUNCTION_NAME "systemd-journal"
21
-#define SYSTEMD_JOURNAL_DEFAULT_TIMEOUT 60
22
-#define SYSTEMD_JOURNAL_MAX_PARAMS 1000
23
-#define SYSTEMD_JOURNAL_DEFAULT_QUERY_DURATION (1 * 3600)
24
-#define SYSTEMD_JOURNAL_DEFAULT_ITEMS_PER_QUERY 200
25
-#define SYSTEMD_JOURNAL_DEFAULT_ITEMS_SAMPLING 1000000
26
-#define SYSTEMD_JOURNAL_SAMPLING_SLOTS 1000
27
-#define SYSTEMD_JOURNAL_SAMPLING_RECALIBRATE 10000
18
+#define SYSTEMD_JOURNAL_SAMPLING_SLOTS 1000
19
+#define SYSTEMD_JOURNAL_SAMPLING_RECALIBRATE 10000
20
29
-#define SYSTEMD_JOURNAL_PROGRESS_EVERY_UT (250 * USEC_PER_MS)
21
+#ifdef HAVE_SD_JOURNAL_RESTART_FIELDS
22
+#define LQS_DEFAULT_SLICE_MODE 1
23
+#else
24
+#define LQS_DEFAULT_SLICE_MODE 0
25
+#endif
26
31
-#define JOURNAL_PARAMETER_HELP "help"
32
-#define JOURNAL_PARAMETER_AFTER "after"
33
-#define JOURNAL_PARAMETER_BEFORE "before"
34
-#define JOURNAL_PARAMETER_ANCHOR "anchor"
35
-#define JOURNAL_PARAMETER_LAST "last"
36
-#define JOURNAL_PARAMETER_QUERY "query"
37
-#define JOURNAL_PARAMETER_FACETS "facets"
38
-#define JOURNAL_PARAMETER_HISTOGRAM "histogram"
39
-#define JOURNAL_PARAMETER_DIRECTION "direction"
40
-#define JOURNAL_PARAMETER_IF_MODIFIED_SINCE "if_modified_since"
41
-#define JOURNAL_PARAMETER_DATA_ONLY "data_only"
42
-#define JOURNAL_PARAMETER_SOURCE "source"
43
-#define JOURNAL_PARAMETER_INFO "info"
44
-#define JOURNAL_PARAMETER_SLICE "slice"
45
-#define JOURNAL_PARAMETER_DELTA "delta"
46
-#define JOURNAL_PARAMETER_TAIL "tail"
47
-#define JOURNAL_PARAMETER_SAMPLING "sampling"
27
+// functions needed by LQS
28
+static SD_JOURNAL_FILE_SOURCE_TYPE get_internal_source_type(const char *value);
29
+
30
+// structures needed by LQS
31
+struct lqs_extension {
32
+ struct {
33
+ usec_t start_ut;
34
+ usec_t stop_ut;
35
+ usec_t first_msg_ut;
36
37
+ sd_id128_t first_msg_writer;
38
+ uint64_t first_msg_seqnum;
39
+ } query_file;
40
+
41
+ struct {
42
+ uint32_t enable_after_samples;
43
+ uint32_t slots;
44
+ uint32_t sampled;
45
+ uint32_t unsampled;
46
+ uint32_t estimated;
47
+ } samples;
48
+
49
+ struct {
50
+ uint32_t enable_after_samples;
51
+ uint32_t every;
52
+ uint32_t skipped;
53
+ uint32_t recalibrate;
54
+ uint32_t sampled;
55
+ uint32_t unsampled;
56
+ uint32_t estimated;
57
+ } samples_per_file;
58
+
59
+ struct {
60
+ usec_t start_ut;
61
+ usec_t end_ut;
62
+ usec_t step_ut;
63
+ uint32_t enable_after_samples;
64
+ uint32_t sampled[SYSTEMD_JOURNAL_SAMPLING_SLOTS];
65
+ uint32_t unsampled[SYSTEMD_JOURNAL_SAMPLING_SLOTS];
66
+ } samples_per_time_slot;
67
+
68
+ // per file progress info
69
+ // size_t cached_count;
70
+
71
+ // progress statistics
72
+ usec_t matches_setup_ut;
73
+ size_t rows_useful;
74
+ size_t rows_read;
75
+ size_t bytes_read;
76
+ size_t files_matched;
77
+ size_t file_working;
78
+};
79
+
80
+// prepare LQS
81
+#define LQS_FUNCTION_NAME SYSTEMD_JOURNAL_FUNCTION_NAME
82
+#define LQS_FUNCTION_DESCRIPTION SYSTEMD_JOURNAL_FUNCTION_DESCRIPTION
83
+#define LQS_DEFAULT_ITEMS_PER_QUERY 200
84
+#define LQS_DEFAULT_ITEMS_SAMPLING 1000000
85
+#define LQS_SOURCE_TYPE SD_JOURNAL_FILE_SOURCE_TYPE
86
+#define LQS_SOURCE_TYPE_ALL SDJF_ALL
87
+#define LQS_SOURCE_TYPE_NONE SDJF_NONE
88
+#define LQS_FUNCTION_GET_INTERNAL_SOURCE_TYPE(value) get_internal_source_type(value)
89
+#define LQS_FUNCTION_SOURCE_TO_JSON_ARRAY(wb) available_journal_file_sources_to_json_array(wb)
90
+#include "libnetdata/facets/logs_query_status.h"
91
+
92
+#include "systemd-journal-sampling.h"
93
+
94
+#define FACET_MAX_VALUE_LENGTH 8192
95
+#define SYSTEMD_JOURNAL_DEFAULT_TIMEOUT 60
96
+#define SYSTEMD_JOURNAL_PROGRESS_EVERY_UT (250 * USEC_PER_MS)
97
#define JOURNAL_KEY_ND_JOURNAL_FILE "ND_JOURNAL_FILE"
98
#define JOURNAL_KEY_ND_JOURNAL_PROCESS "ND_JOURNAL_PROCESS"
51
-
52
-#define JOURNAL_DEFAULT_SLICE_MODE true
99
#define JOURNAL_DEFAULT_DIRECTION FACETS_ANCHOR_DIRECTION_BACKWARD
54
-
100
#define SYSTEMD_ALWAYS_VISIBLE_KEYS NULL
101
102
#define SYSTEMD_KEYS_EXCLUDED_FROM_FACETS \
227
228
// ----------------------------------------------------------------------------
229
185
-typedef struct function_query_status {
186
- bool *cancelled; // a pointer to the cancelling boolean
187
- usec_t *stop_monotonic_ut;
188
-
189
- // request
190
- const char *transaction;
191
-
192
- SD_JOURNAL_FILE_SOURCE_TYPE source_type;
193
- SIMPLE_PATTERN *sources;
194
- usec_t after_ut;
195
- usec_t before_ut;
196
-
197
- struct {
198
- usec_t start_ut;
199
- usec_t stop_ut;
200
- } anchor;
201
-
202
- FACETS_ANCHOR_DIRECTION direction;
203
- size_t entries;
204
- usec_t if_modified_since;
205
- bool delta;
206
- bool tail;
207
- bool data_only;
208
- bool slice;
209
- size_t sampling;
210
- size_t filters;
211
- usec_t last_modified;
212
- const char *query;
213
- const char *histogram;
214
-
215
- struct {
216
- usec_t start_ut; // the starting time of the query - we start from this
217
- usec_t stop_ut; // the ending time of the query - we stop at this
218
- usec_t first_msg_ut;
219
-
220
- sd_id128_t first_msg_writer;
221
- uint64_t first_msg_seqnum;
222
- } query_file;
223
-
224
- struct {
225
- uint32_t enable_after_samples;
226
- uint32_t slots;
227
- uint32_t sampled;
228
- uint32_t unsampled;
229
- uint32_t estimated;
230
- } samples;
231
-
232
- struct {
233
- uint32_t enable_after_samples;
234
- uint32_t every;
235
- uint32_t skipped;
236
- uint32_t recalibrate;
237
- uint32_t sampled;
238
- uint32_t unsampled;
239
- uint32_t estimated;
240
- } samples_per_file;
241
-
242
- struct {
243
- usec_t start_ut;
244
- usec_t end_ut;
245
- usec_t step_ut;
246
- uint32_t enable_after_samples;
247
- uint32_t sampled[SYSTEMD_JOURNAL_SAMPLING_SLOTS];
248
- uint32_t unsampled[SYSTEMD_JOURNAL_SAMPLING_SLOTS];
249
- } samples_per_time_slot;
250
-
251
- // per file progress info
252
- // size_t cached_count;
230
+static SD_JOURNAL_FILE_SOURCE_TYPE get_internal_source_type(const char *value) {
231
+ if(strcmp(value, SDJF_SOURCE_ALL_NAME) == 0)
232
+ return SDJF_ALL;
233
+ else if(strcmp(value, SDJF_SOURCE_LOCAL_NAME) == 0)
234
+ return SDJF_LOCAL_ALL;
235
+ else if(strcmp(value, SDJF_SOURCE_REMOTES_NAME) == 0)
236
+ return SDJF_REMOTE_ALL;
237
+ else if(strcmp(value, SDJF_SOURCE_NAMESPACES_NAME) == 0)
238
+ return SDJF_LOCAL_NAMESPACE;
239
+ else if(strcmp(value, SDJF_SOURCE_LOCAL_SYSTEM_NAME) == 0)
240
+ return SDJF_LOCAL_SYSTEM;
241
+ else if(strcmp(value, SDJF_SOURCE_LOCAL_USERS_NAME) == 0)
242
+ return SDJF_LOCAL_USER;
243
+ else if(strcmp(value, SDJF_SOURCE_LOCAL_OTHER_NAME) == 0)
244
+ return SDJF_LOCAL_OTHER;
245
254
- // progress statistics
255
- usec_t matches_setup_ut;
256
- size_t rows_useful;
257
- size_t rows_read;
258
- size_t bytes_read;
259
- size_t files_matched;
260
- size_t file_working;
261
-} FUNCTION_QUERY_STATUS;
262
-
263
-static void log_fqs(FUNCTION_QUERY_STATUS *fqs, const char *msg) {
264
- netdata_log_error("ERROR: %s, on query "
265
- "timeframe [%"PRIu64" - %"PRIu64"], "
266
- "anchor [%"PRIu64" - %"PRIu64"], "
267
- "if_modified_since %"PRIu64", "
268
- "data_only:%s, delta:%s, tail:%s, direction:%s"
269
- , msg
270
- , fqs->after_ut, fqs->before_ut
271
- , fqs->anchor.start_ut, fqs->anchor.stop_ut
272
- , fqs->if_modified_since
273
- , fqs->data_only ? "true" : "false"
274
- , fqs->delta ? "true" : "false"
275
- , fqs->tail ? "tail" : "false"
276
- , fqs->direction == FACETS_ANCHOR_DIRECTION_FORWARD ? "forward" : "backward");
246
+ return SDJF_NONE;
247
}
248
249
+// ----------------------------------------------------------------------------
250
+
251
static inline bool netdata_systemd_journal_seek_to(sd_journal *j, usec_t timestamp) {
252
if(sd_journal_seek_realtime_usec(j, timestamp) < 0) {
253
netdata_log_error("SYSTEMD-JOURNAL: Failed to seek to %" PRIu64, timestamp);
262
263
#define JD_SOURCE_REALTIME_TIMESTAMP "_SOURCE_REALTIME_TIMESTAMP"
264
293
-// ----------------------------------------------------------------------------
294
-// sampling support
295
-
296
-static void sampling_query_init(FUNCTION_QUERY_STATUS *fqs, FACETS *facets) {
297
- if(!fqs->sampling)
298
- return;
299
-
300
- if(!fqs->slice) {
301
- // the user is doing a full data query
302
- // disable sampling
303
- fqs->sampling = 0;
304
- return;
305
- }
306
-
307
- if(fqs->data_only) {
308
- // the user is doing a data query
309
- // disable sampling
310
- fqs->sampling = 0;
311
- return;
312
- }
313
-
314
- if(!fqs->files_matched) {
315
- // no files have been matched
316
- // disable sampling
317
- fqs->sampling = 0;
318
- return;
319
- }
320
-
321
- fqs->samples.slots = facets_histogram_slots(facets);
322
- if(fqs->samples.slots < 2) fqs->samples.slots = 2;
323
- if(fqs->samples.slots > SYSTEMD_JOURNAL_SAMPLING_SLOTS)
324
- fqs->samples.slots = SYSTEMD_JOURNAL_SAMPLING_SLOTS;
325
-
326
- if(!fqs->after_ut || !fqs->before_ut || fqs->after_ut >= fqs->before_ut) {
327
- // we don't have enough information for sampling
328
- fqs->sampling = 0;
329
- return;
330
- }
331
-
332
- usec_t delta = fqs->before_ut - fqs->after_ut;
333
- usec_t step = delta / facets_histogram_slots(facets) - 1;
334
- if(step < 1) step = 1;
335
-
336
- fqs->samples_per_time_slot.start_ut = fqs->after_ut;
337
- fqs->samples_per_time_slot.end_ut = fqs->before_ut;
338
- fqs->samples_per_time_slot.step_ut = step;
339
-
340
- // the minimum number of rows to enable sampling
341
- fqs->samples.enable_after_samples = fqs->sampling / 2;
342
-
343
- size_t files_matched = fqs->files_matched;
344
- if(!files_matched)
345
- files_matched = 1;
346
-
347
- // the minimum number of rows per file to enable sampling
348
- fqs->samples_per_file.enable_after_samples = (fqs->sampling / 4) / files_matched;
349
- if(fqs->samples_per_file.enable_after_samples < fqs->entries)
350
- fqs->samples_per_file.enable_after_samples = fqs->entries;
351
-
352
- // the minimum number of rows per time slot to enable sampling
353
- fqs->samples_per_time_slot.enable_after_samples = (fqs->sampling / 4) / fqs->samples.slots;
354
- if(fqs->samples_per_time_slot.enable_after_samples < fqs->entries)
355
- fqs->samples_per_time_slot.enable_after_samples = fqs->entries;
356
-}
357
-
358
-static void sampling_file_init(FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf __maybe_unused) {
359
- fqs->samples_per_file.sampled = 0;
360
- fqs->samples_per_file.unsampled = 0;
361
- fqs->samples_per_file.estimated = 0;
362
- fqs->samples_per_file.every = 0;
363
- fqs->samples_per_file.skipped = 0;
364
- fqs->samples_per_file.recalibrate = 0;
365
-}
366
-
367
-static size_t sampling_file_lines_scanned_so_far(FUNCTION_QUERY_STATUS *fqs) {
368
- size_t sampled = fqs->samples_per_file.sampled + fqs->samples_per_file.unsampled;
369
- if(!sampled) sampled = 1;
370
- return sampled;
371
-}
372
-
373
-static void sampling_running_file_query_overlapping_timeframe_ut(
374
- FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf, FACETS_ANCHOR_DIRECTION direction,
375
- usec_t msg_ut, usec_t *after_ut, usec_t *before_ut) {
376
-
377
- // find the overlap of the query and file timeframes
378
- // taking into account the first message we encountered
379
-
380
- usec_t oldest_ut, newest_ut;
381
- if(direction == FACETS_ANCHOR_DIRECTION_FORWARD) {
382
- // the first message we know (oldest)
383
- oldest_ut = fqs->query_file.first_msg_ut ? fqs->query_file.first_msg_ut : jf->msg_first_ut;
384
- if(!oldest_ut) oldest_ut = fqs->query_file.start_ut;
385
-
386
- if(jf->msg_last_ut)
387
- newest_ut = MIN(fqs->query_file.stop_ut, jf->msg_last_ut);
388
- else if(jf->file_last_modified_ut)
389
- newest_ut = MIN(fqs->query_file.stop_ut, jf->file_last_modified_ut);
390
- else
391
- newest_ut = fqs->query_file.stop_ut;
392
-
393
- if(msg_ut < oldest_ut)
394
- oldest_ut = msg_ut - 1;
395
- }
396
- else /* BACKWARD */ {
397
- // the latest message we know (newest)
398
- newest_ut = fqs->query_file.first_msg_ut ? fqs->query_file.first_msg_ut : jf->msg_last_ut;
399
- if(!newest_ut) newest_ut = fqs->query_file.start_ut;
400
-
401
- if(jf->msg_first_ut)
402
- oldest_ut = MAX(fqs->query_file.stop_ut, jf->msg_first_ut);
403
- else
404
- oldest_ut = fqs->query_file.stop_ut;
405
-
406
- if(newest_ut < msg_ut)
407
- newest_ut = msg_ut + 1;
408
- }
409
-
410
- *after_ut = oldest_ut;
411
- *before_ut = newest_ut;
412
-}
413
-
414
-static double sampling_running_file_query_progress_by_time(FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf,
415
- FACETS_ANCHOR_DIRECTION direction, usec_t msg_ut) {
416
-
417
- usec_t after_ut, before_ut, elapsed_ut;
418
- sampling_running_file_query_overlapping_timeframe_ut(fqs, jf, direction, msg_ut, &after_ut, &before_ut);
419
-
420
- if(direction == FACETS_ANCHOR_DIRECTION_FORWARD)
421
- elapsed_ut = msg_ut - after_ut;
422
- else
423
- elapsed_ut = before_ut - msg_ut;
424
-
425
- usec_t total_ut = before_ut - after_ut;
426
- double progress = (double)elapsed_ut / (double)total_ut;
427
-
428
- return progress;
429
-}
430
-
431
-static usec_t sampling_running_file_query_remaining_time(FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf,
432
- FACETS_ANCHOR_DIRECTION direction, usec_t msg_ut,
433
- usec_t *total_time_ut, usec_t *remaining_start_ut,
434
- usec_t *remaining_end_ut) {
435
- usec_t after_ut, before_ut;
436
- sampling_running_file_query_overlapping_timeframe_ut(fqs, jf, direction, msg_ut, &after_ut, &before_ut);
437
-
438
- // since we have a timestamp in msg_ut
439
- // this timestamp can extend the overlap
440
- if(msg_ut <= after_ut)
441
- after_ut = msg_ut - 1;
442
-
443
- if(msg_ut >= before_ut)
444
- before_ut = msg_ut + 1;
445
-
446
- // return the remaining duration
447
- usec_t remaining_from_ut, remaining_to_ut;
448
- if(direction == FACETS_ANCHOR_DIRECTION_FORWARD) {
449
- remaining_from_ut = msg_ut;
450
- remaining_to_ut = before_ut;
451
- }
452
- else {
453
- remaining_from_ut = after_ut;
454
- remaining_to_ut = msg_ut;
455
- }
456
-
457
- usec_t remaining_ut = remaining_to_ut - remaining_from_ut;
458
-
459
- if(total_time_ut)
460
- *total_time_ut = (before_ut > after_ut) ? before_ut - after_ut : 1;
461
-
462
- if(remaining_start_ut)
463
- *remaining_start_ut = remaining_from_ut;
464
-
465
- if(remaining_end_ut)
466
- *remaining_end_ut = remaining_to_ut;
467
-
468
- return remaining_ut;
469
-}
470
-
471
-static size_t sampling_running_file_query_estimate_remaining_lines_by_time(FUNCTION_QUERY_STATUS *fqs,
472
- struct journal_file *jf,
473
- FACETS_ANCHOR_DIRECTION direction,
474
- usec_t msg_ut) {
475
- size_t scanned_lines = sampling_file_lines_scanned_so_far(fqs);
476
-
477
- // Calculate the proportion of time covered
478
- usec_t total_time_ut, remaining_start_ut, remaining_end_ut;
479
- usec_t remaining_time_ut = sampling_running_file_query_remaining_time(fqs, jf, direction, msg_ut, &total_time_ut,
480
- &remaining_start_ut, &remaining_end_ut);
481
- if (total_time_ut == 0) total_time_ut = 1;
482
-
483
- double proportion_by_time = (double) (total_time_ut - remaining_time_ut) / (double) total_time_ut;
484
-
485
- if (proportion_by_time == 0 || proportion_by_time > 1.0 || !isfinite(proportion_by_time))
486
- proportion_by_time = 1.0;
487
-
488
- // Estimate the total number of lines in the file
489
- size_t expected_matching_logs_by_time = (size_t)((double)scanned_lines / proportion_by_time);
490
-
491
- if(jf->messages_in_file && expected_matching_logs_by_time > jf->messages_in_file)
492
- expected_matching_logs_by_time = jf->messages_in_file;
493
-
494
- // Calculate the estimated number of remaining lines
495
- size_t remaining_logs_by_time = expected_matching_logs_by_time - scanned_lines;
496
- if (remaining_logs_by_time < 1) remaining_logs_by_time = 1;
497
-
498
-// nd_log(NDLS_COLLECTORS, NDLP_INFO,
499
-// "JOURNAL ESTIMATION: '%s' "
500
-// "scanned_lines=%zu [sampled=%zu, unsampled=%zu, estimated=%zu], "
501
-// "file [%"PRIu64" - %"PRIu64", duration %"PRId64", known lines in file %zu], "
502
-// "query [%"PRIu64" - %"PRIu64", duration %"PRId64"], "
503
-// "first message read from the file at %"PRIu64", current message at %"PRIu64", "
504
-// "proportion of time %.2f %%, "
505
-// "expected total lines in file %zu, "
506
-// "remaining lines %zu, "
507
-// "remaining time %"PRIu64" [%"PRIu64" - %"PRIu64", duration %"PRId64"]"
508
-// , jf->filename
509
-// , scanned_lines, fqs->samples_per_file.sampled, fqs->samples_per_file.unsampled, fqs->samples_per_file.estimated
510
-// , jf->msg_first_ut, jf->msg_last_ut, jf->msg_last_ut - jf->msg_first_ut, jf->messages_in_file
511
-// , fqs->query_file.start_ut, fqs->query_file.stop_ut, fqs->query_file.stop_ut - fqs->query_file.start_ut
512
-// , fqs->query_file.first_msg_ut, msg_ut
513
-// , proportion_by_time * 100.0
514
-// , expected_matching_logs_by_time
515
-// , remaining_logs_by_time
516
-// , remaining_time_ut, remaining_start_ut, remaining_end_ut, remaining_end_ut - remaining_start_ut
517
-// );
518
-
519
- return remaining_logs_by_time;
520
-}
521
-
522
-static size_t sampling_running_file_query_estimate_remaining_lines(sd_journal *j __maybe_unused, FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf, FACETS_ANCHOR_DIRECTION direction, usec_t msg_ut) {
523
- size_t remaining_logs_by_seqnum = 0;
524
-
525
-#ifdef HAVE_SD_JOURNAL_GET_SEQNUM
526
- size_t expected_matching_logs_by_seqnum = 0;
527
- double proportion_by_seqnum = 0.0;
528
- uint64_t current_msg_seqnum;
529
- sd_id128_t current_msg_writer;
530
- if(!fqs->query_file.first_msg_seqnum || sd_journal_get_seqnum(j, ¤t_msg_seqnum, ¤t_msg_writer) < 0) {
531
- fqs->query_file.first_msg_seqnum = 0;
532
- fqs->query_file.first_msg_writer = SD_ID128_NULL;
533
- }
534
- else if(jf->messages_in_file) {
535
- size_t scanned_lines = sampling_file_lines_scanned_so_far(fqs);
536
-
537
- double proportion_of_all_lines_so_far;
538
- if(direction == FACETS_ANCHOR_DIRECTION_FORWARD)
539
- proportion_of_all_lines_so_far = (double)scanned_lines / (double)(current_msg_seqnum - jf->first_seqnum);
540
- else
541
- proportion_of_all_lines_so_far = (double)scanned_lines / (double)(jf->last_seqnum - current_msg_seqnum);
542
-
543
- if(proportion_of_all_lines_so_far > 1.0)
544
- proportion_of_all_lines_so_far = 1.0;
545
-
546
- expected_matching_logs_by_seqnum = (size_t)(proportion_of_all_lines_so_far * (double)jf->messages_in_file);
547
-
548
- proportion_by_seqnum = (double)scanned_lines / (double)expected_matching_logs_by_seqnum;
549
-
550
- if (proportion_by_seqnum == 0 || proportion_by_seqnum > 1.0 || !isfinite(proportion_by_seqnum))
551
- proportion_by_seqnum = 1.0;
552
-
553
- remaining_logs_by_seqnum = expected_matching_logs_by_seqnum - scanned_lines;
554
- if(!remaining_logs_by_seqnum) remaining_logs_by_seqnum = 1;
555
- }
556
-#endif
557
-
558
- if(remaining_logs_by_seqnum)
559
- return remaining_logs_by_seqnum;
560
-
561
- return sampling_running_file_query_estimate_remaining_lines_by_time(fqs, jf, direction, msg_ut);
562
-}
563
-
564
-static void sampling_decide_file_sampling_every(sd_journal *j, FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf, FACETS_ANCHOR_DIRECTION direction, usec_t msg_ut) {
565
- size_t files_matched = fqs->files_matched;
566
- if(!files_matched) files_matched = 1;
567
-
568
- size_t remaining_lines = sampling_running_file_query_estimate_remaining_lines(j, fqs, jf, direction, msg_ut);
569
- size_t wanted_samples = (fqs->sampling / 2) / files_matched;
570
- if(!wanted_samples) wanted_samples = 1;
571
-
572
- fqs->samples_per_file.every = remaining_lines / wanted_samples;
573
-
574
- if(fqs->samples_per_file.every < 1)
575
- fqs->samples_per_file.every = 1;
576
-}
577
-
578
-typedef enum {
579
- SAMPLING_STOP_AND_ESTIMATE = -1,
580
- SAMPLING_FULL = 0,
581
- SAMPLING_SKIP_FIELDS = 1,
582
-} sampling_t;
583
-
584
-static inline sampling_t is_row_in_sample(sd_journal *j, FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf, usec_t msg_ut, FACETS_ANCHOR_DIRECTION direction, bool candidate_to_keep) {
585
- if(!fqs->sampling || candidate_to_keep)
586
- return SAMPLING_FULL;
587
-
588
- if(unlikely(msg_ut < fqs->samples_per_time_slot.start_ut))
589
- msg_ut = fqs->samples_per_time_slot.start_ut;
590
- if(unlikely(msg_ut > fqs->samples_per_time_slot.end_ut))
591
- msg_ut = fqs->samples_per_time_slot.end_ut;
592
-
593
- size_t slot = (msg_ut - fqs->samples_per_time_slot.start_ut) / fqs->samples_per_time_slot.step_ut;
594
- if(slot >= fqs->samples.slots)
595
- slot = fqs->samples.slots - 1;
596
-
597
- bool should_sample = false;
598
-
599
- if(fqs->samples.sampled < fqs->samples.enable_after_samples ||
600
- fqs->samples_per_file.sampled < fqs->samples_per_file.enable_after_samples ||
601
- fqs->samples_per_time_slot.sampled[slot] < fqs->samples_per_time_slot.enable_after_samples)
602
- should_sample = true;
603
-
604
- else if(fqs->samples_per_file.recalibrate >= SYSTEMD_JOURNAL_SAMPLING_RECALIBRATE || !fqs->samples_per_file.every) {
605
- // this is the first to be unsampled for this file
606
- sampling_decide_file_sampling_every(j, fqs, jf, direction, msg_ut);
607
- fqs->samples_per_file.recalibrate = 0;
608
- should_sample = true;
609
- }
610
- else {
611
- // we sample 1 every fqs->samples_per_file.every
612
- if(fqs->samples_per_file.skipped >= fqs->samples_per_file.every) {
613
- fqs->samples_per_file.skipped = 0;
614
- should_sample = true;
615
- }
616
- else
617
- fqs->samples_per_file.skipped++;
618
- }
619
-
620
- if(should_sample) {
621
- fqs->samples.sampled++;
622
- fqs->samples_per_file.sampled++;
623
- fqs->samples_per_time_slot.sampled[slot]++;
624
-
625
- return SAMPLING_FULL;
626
- }
627
-
628
- fqs->samples_per_file.recalibrate++;
629
-
630
- fqs->samples.unsampled++;
631
- fqs->samples_per_file.unsampled++;
632
- fqs->samples_per_time_slot.unsampled[slot]++;
633
-
634
- if(fqs->samples_per_file.unsampled > fqs->samples_per_file.sampled) {
635
- double progress_by_time = sampling_running_file_query_progress_by_time(fqs, jf, direction, msg_ut);
636
-
637
- if(progress_by_time > SYSTEMD_JOURNAL_ENABLE_ESTIMATIONS_FILE_PERCENTAGE)
638
- return SAMPLING_STOP_AND_ESTIMATE;
639
- }
640
-
641
- return SAMPLING_SKIP_FIELDS;
642
-}
643
-
644
-static void sampling_update_running_query_file_estimates(FACETS *facets, sd_journal *j, FUNCTION_QUERY_STATUS *fqs, struct journal_file *jf, usec_t msg_ut, FACETS_ANCHOR_DIRECTION direction) {
645
- usec_t total_time_ut, remaining_start_ut, remaining_end_ut;
646
- sampling_running_file_query_remaining_time(fqs, jf, direction, msg_ut, &total_time_ut, &remaining_start_ut,
647
- &remaining_end_ut);
648
- size_t remaining_lines = sampling_running_file_query_estimate_remaining_lines(j, fqs, jf, direction, msg_ut);
649
- facets_update_estimations(facets, remaining_start_ut, remaining_end_ut, remaining_lines);
650
- fqs->samples.estimated += remaining_lines;
651
- fqs->samples_per_file.estimated += remaining_lines;
652
-}
653
-
265
// ----------------------------------------------------------------------------
266
267
static inline size_t netdata_systemd_journal_process_row(sd_journal *j, FACETS *facets, struct journal_file *jf, usec_t *msg_ut) {
332
333
ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_backward(
334
sd_journal *j, BUFFER *wb __maybe_unused, FACETS *facets,
724
- struct journal_file *jf, FUNCTION_QUERY_STATUS *fqs) {
335
+ struct journal_file *jf,
336
+ LOGS_QUERY_STATUS *fqs) {
337
338
usec_t anchor_delta = __atomic_load_n(&jf->max_journal_vs_realtime_delta_ut, __ATOMIC_RELAXED);
339
728
- usec_t start_ut = ((fqs->data_only && fqs->anchor.start_ut) ? fqs->anchor.start_ut : fqs->before_ut) + anchor_delta;
729
- usec_t stop_ut = (fqs->data_only && fqs->anchor.stop_ut) ? fqs->anchor.stop_ut : fqs->after_ut;
730
- bool stop_when_full = (fqs->data_only && !fqs->anchor.stop_ut);
340
+ usec_t start_ut = ((fqs->rq.data_only && fqs->anchor.start_ut) ? fqs->anchor.start_ut : fqs->rq.before_ut) + anchor_delta;
341
+ usec_t stop_ut = (fqs->rq.data_only && fqs->anchor.stop_ut) ? fqs->anchor.stop_ut : fqs->rq.after_ut;
342
+ bool stop_when_full = (fqs->rq.data_only && !fqs->anchor.stop_ut);
343
732
- fqs->query_file.start_ut = start_ut;
733
- fqs->query_file.stop_ut = stop_ut;
344
+ fqs->c.query_file.start_ut = start_ut;
345
+ fqs->c.query_file.stop_ut = stop_ut;
346
347
if(!netdata_systemd_journal_seek_to(j, start_ut))
348
return ND_SD_JOURNAL_FAILED_TO_SEEK;
377
378
if(unlikely(!first_msg_ut)) {
379
first_msg_ut = msg_ut;
768
- fqs->query_file.first_msg_ut = msg_ut;
380
+ fqs->c.query_file.first_msg_ut = msg_ut;
381
382
#ifdef HAVE_SD_JOURNAL_GET_SEQNUM
771
- if(sd_journal_get_seqnum(j, &fqs->query_file.first_msg_seqnum, &fqs->query_file.first_msg_writer) < 0) {
772
- fqs->query_file.first_msg_seqnum = 0;
773
- fqs->query_file.first_msg_writer = SD_ID128_NULL;
383
+ if(sd_journal_get_seqnum(j, &fqs->c.query_file.first_msg_seqnum, &fqs->c.query_file.first_msg_writer) < 0) {
384
+ fqs->c.query_file.first_msg_seqnum = 0;
385
+ fqs->c.query_file.first_msg_writer = SD_ID128_NULL;
386
}
387
#endif
388
}
406
row_counter++;
407
if(unlikely((row_counter % FUNCTION_DATA_ONLY_CHECK_EVERY_ROWS) == 0 &&
408
stop_when_full &&
797
- facets_rows(facets) >= fqs->entries)) {
409
+ facets_rows(facets) >= fqs->rq.entries)) {
410
// stop the data only query
411
usec_t oldest = facets_row_oldest_ut(facets);
412
if(oldest && msg_ut < (oldest - anchor_delta))
414
}
415
416
if(unlikely(row_counter % FUNCTION_PROGRESS_EVERY_ROWS == 0)) {
805
- FUNCTION_PROGRESS_UPDATE_ROWS(fqs->rows_read, row_counter - last_row_counter);
417
+ FUNCTION_PROGRESS_UPDATE_ROWS(fqs->c.rows_read, row_counter - last_row_counter);
418
last_row_counter = row_counter;
419
808
- FUNCTION_PROGRESS_UPDATE_BYTES(fqs->bytes_read, bytes - last_bytes);
420
+ FUNCTION_PROGRESS_UPDATE_BYTES(fqs->c.bytes_read, bytes - last_bytes);
421
last_bytes = bytes;
422
423
status = check_stop(fqs->cancelled, fqs->stop_monotonic_ut);
431
}
432
}
433
822
- FUNCTION_PROGRESS_UPDATE_ROWS(fqs->rows_read, row_counter - last_row_counter);
823
- FUNCTION_PROGRESS_UPDATE_BYTES(fqs->bytes_read, bytes - last_bytes);
434
+ FUNCTION_PROGRESS_UPDATE_ROWS(fqs->c.rows_read, row_counter - last_row_counter);
435
+ FUNCTION_PROGRESS_UPDATE_BYTES(fqs->c.bytes_read, bytes - last_bytes);
436
825
- fqs->rows_useful += rows_useful;
437
+ fqs->c.rows_useful += rows_useful;
438
439
if(errors_no_timestamp)
440
netdata_log_error("SYSTEMD-JOURNAL: %zu lines did not have timestamps", errors_no_timestamp);
447
448
ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_forward(
449
sd_journal *j, BUFFER *wb __maybe_unused, FACETS *facets,
838
- struct journal_file *jf, FUNCTION_QUERY_STATUS *fqs) {
450
+ struct journal_file *jf,
451
+ LOGS_QUERY_STATUS *fqs) {
452
453
usec_t anchor_delta = __atomic_load_n(&jf->max_journal_vs_realtime_delta_ut, __ATOMIC_RELAXED);
454
842
- usec_t start_ut = (fqs->data_only && fqs->anchor.start_ut) ? fqs->anchor.start_ut : fqs->after_ut;
843
- usec_t stop_ut = ((fqs->data_only && fqs->anchor.stop_ut) ? fqs->anchor.stop_ut : fqs->before_ut) + anchor_delta;
844
- bool stop_when_full = (fqs->data_only && !fqs->anchor.stop_ut);
455
+ usec_t start_ut = (fqs->rq.data_only && fqs->anchor.start_ut) ? fqs->anchor.start_ut : fqs->rq.after_ut;
456
+ usec_t stop_ut = ((fqs->rq.data_only && fqs->anchor.stop_ut) ? fqs->anchor.stop_ut : fqs->rq.before_ut) + anchor_delta;
457
+ bool stop_when_full = (fqs->rq.data_only && !fqs->anchor.stop_ut);
458
846
- fqs->query_file.start_ut = start_ut;
847
- fqs->query_file.stop_ut = stop_ut;
459
+ fqs->c.query_file.start_ut = start_ut;
460
+ fqs->c.query_file.stop_ut = stop_ut;
461
462
if(!netdata_systemd_journal_seek_to(j, start_ut))
463
return ND_SD_JOURNAL_FAILED_TO_SEEK;
492
493
if(unlikely(!first_msg_ut)) {
494
first_msg_ut = msg_ut;
882
- fqs->query_file.first_msg_ut = msg_ut;
495
+ fqs->c.query_file.first_msg_ut = msg_ut;
496
}
497
498
sampling_t sample = is_row_in_sample(j, fqs, jf, msg_ut,
514
row_counter++;
515
if(unlikely((row_counter % FUNCTION_DATA_ONLY_CHECK_EVERY_ROWS) == 0 &&
516
stop_when_full &&
904
- facets_rows(facets) >= fqs->entries)) {
517
+ facets_rows(facets) >= fqs->rq.entries)) {
518
// stop the data only query
519
usec_t newest = facets_row_newest_ut(facets);
520
if(newest && msg_ut > (newest + anchor_delta))
522
}
523
524
if(unlikely(row_counter % FUNCTION_PROGRESS_EVERY_ROWS == 0)) {
912
- FUNCTION_PROGRESS_UPDATE_ROWS(fqs->rows_read, row_counter - last_row_counter);
525
+ FUNCTION_PROGRESS_UPDATE_ROWS(fqs->c.rows_read, row_counter - last_row_counter);
526
last_row_counter = row_counter;
527
915
- FUNCTION_PROGRESS_UPDATE_BYTES(fqs->bytes_read, bytes - last_bytes);
528
+ FUNCTION_PROGRESS_UPDATE_BYTES(fqs->c.bytes_read, bytes - last_bytes);
529
last_bytes = bytes;
530
531
status = check_stop(fqs->cancelled, fqs->stop_monotonic_ut);
539
}
540
}
541
929
- FUNCTION_PROGRESS_UPDATE_ROWS(fqs->rows_read, row_counter - last_row_counter);
930
- FUNCTION_PROGRESS_UPDATE_BYTES(fqs->bytes_read, bytes - last_bytes);
542
+ FUNCTION_PROGRESS_UPDATE_ROWS(fqs->c.rows_read, row_counter - last_row_counter);
543
+ FUNCTION_PROGRESS_UPDATE_BYTES(fqs->c.bytes_read, bytes - last_bytes);
544
932
- fqs->rows_useful += rows_useful;
545
+ fqs->c.rows_useful += rows_useful;
546
547
if(errors_no_timestamp)
548
netdata_log_error("SYSTEMD-JOURNAL: %zu lines did not have timestamps", errors_no_timestamp);
576
}
577
578
#ifdef HAVE_SD_JOURNAL_RESTART_FIELDS
966
-static bool netdata_systemd_filtering_by_journal(sd_journal *j, FACETS *facets, FUNCTION_QUERY_STATUS *fqs) {
579
+static bool netdata_systemd_filtering_by_journal(sd_journal *j, FACETS *facets, LOGS_QUERY_STATUS *lqs) {
580
const char *field = NULL;
581
const void *data = NULL;
582
size_t data_length;
587
SD_JOURNAL_FOREACH_FIELD(j, field) { // for each key
588
bool interesting;
589
977
- if(fqs->data_only)
590
+ if(lqs->rq.data_only)
591
interesting = facets_key_name_is_filter(facets, field);
592
else
593
interesting = facets_key_name_is_facet(facets, field);
636
}
637
638
if(failures) {
1026
- log_fqs(fqs, "failed to setup journal filter, will run the full query.");
639
+ lqs_log_error(lqs, "failed to setup journal filter, will run the full query.");
640
sd_journal_flush_matches(j);
641
return true;
642
}
647
648
static ND_SD_JOURNAL_STATUS netdata_systemd_journal_query_one_file(
649
const char *filename, BUFFER *wb, FACETS *facets,
1037
- struct journal_file *jf, FUNCTION_QUERY_STATUS *fqs) {
650
+ struct journal_file *jf,
651
+ LOGS_QUERY_STATUS *fqs) {
652
653
sd_journal *j = NULL;
654
errno_clear();
670
bool matches_filters = true;
671
672
#ifdef HAVE_SD_JOURNAL_RESTART_FIELDS
1059
- if(fqs->slice) {
673
+ if(fqs->rq.slice) {
674
usec_t started = now_monotonic_usec();
675
1062
- matches_filters = netdata_systemd_filtering_by_journal(j, facets, fqs) || !fqs->filters;
676
+ matches_filters = netdata_systemd_filtering_by_journal(j, facets, fqs) || !fqs->rq.filters;
677
usec_t ended = now_monotonic_usec();
678
1065
- fqs->matches_setup_ut += (ended - started);
679
+ fqs->c.matches_setup_ut += (ended - started);
680
}
681
#endif // HAVE_SD_JOURNAL_RESTART_FIELDS
682
683
if(matches_filters) {
1070
- if(fqs->direction == FACETS_ANCHOR_DIRECTION_FORWARD)
684
+ if(fqs->rq.direction == FACETS_ANCHOR_DIRECTION_FORWARD)
685
status = netdata_systemd_journal_query_forward(j, wb, facets, jf, fqs);
686
else
687
status = netdata_systemd_journal_query_backward(j, wb, facets, jf, fqs);
695
return status;
696
}
697
1084
-static bool jf_is_mine(struct journal_file *jf, FUNCTION_QUERY_STATUS *fqs) {
698
+static bool jf_is_mine(struct journal_file *jf, LOGS_QUERY_STATUS *fqs) {
699
1086
- if((fqs->source_type == SDJF_NONE && !fqs->sources) || (jf->source_type & fqs->source_type) ||
1087
- (fqs->sources && simple_pattern_matches(fqs->sources, string2str(jf->source)))) {
700
+ if((fqs->rq.source_type == SDJF_NONE && !fqs->rq.sources) || (jf->source_type & fqs->rq.source_type) ||
701
+ (fqs->rq.sources && simple_pattern_matches(fqs->rq.sources, string2str(jf->source)))) {
702
703
if(!jf->msg_last_ut)
704
// the file is not scanned yet, or the timestamps have not been updated,
709
usec_t first_ut = jf->msg_first_ut - anchor_delta;
710
usec_t last_ut = jf->msg_last_ut + anchor_delta;
711
1098
- if(last_ut >= fqs->after_ut && first_ut <= fqs->before_ut)
712
+ if(last_ut >= fqs->rq.after_ut && first_ut <= fqs->rq.before_ut)
713
return true;
714
}
715
716
return false;
717
}
718
1105
-static int netdata_systemd_journal_query(BUFFER *wb, FACETS *facets, FUNCTION_QUERY_STATUS *fqs) {
719
+static int netdata_systemd_journal_query(BUFFER *wb, LOGS_QUERY_STATUS *lqs) {
720
+ FACETS *facets = lqs->facets;
721
+
722
ND_SD_JOURNAL_STATUS status = ND_SD_JOURNAL_NO_FILE_MATCHED;
723
struct journal_file *jf;
724
1109
- fqs->files_matched = 0;
1110
- fqs->file_working = 0;
1111
- fqs->rows_useful = 0;
1112
- fqs->rows_read = 0;
1113
- fqs->bytes_read = 0;
725
+ lqs->c.files_matched = 0;
726
+ lqs->c.file_working = 0;
727
+ lqs->c.rows_useful = 0;
728
+ lqs->c.rows_read = 0;
729
+ lqs->c.bytes_read = 0;
730
731
size_t files_used = 0;
732
size_t files_max = dictionary_entries(journal_files_registry);
735
// count the files
736
bool files_are_newer = false;
737
dfe_start_read(journal_files_registry, jf) {
1122
- if(!jf_is_mine(jf, fqs))
738
+ if(!jf_is_mine(jf, lqs))
739
continue;
740
741
file_items[files_used++] = dictionary_acquired_item_dup(journal_files_registry, jf_dfe.item);
742
1127
- if(jf->msg_last_ut > fqs->if_modified_since)
743
+ if(jf->msg_last_ut > lqs->rq.if_modified_since)
744
files_are_newer = true;
745
}
746
dfe_done(jf);
747
1132
- fqs->files_matched = files_used;
748
+ lqs->c.files_matched = files_used;
749
1134
- if(fqs->if_modified_since && !files_are_newer) {
1135
- buffer_flush(wb);
1136
- return HTTP_RESP_NOT_MODIFIED;
1137
- }
750
+ if(lqs->rq.if_modified_since && !files_are_newer)
751
+ return rrd_call_function_error(wb, "not modified", HTTP_RESP_NOT_MODIFIED);
752
753
// sort the files, so that they are optimal for facets
754
if(files_used >= 2) {
1141
- if (fqs->direction == FACETS_ANCHOR_DIRECTION_BACKWARD)
755
+ if (lqs->rq.direction == FACETS_ANCHOR_DIRECTION_BACKWARD)
756
qsort(file_items, files_used, sizeof(const DICTIONARY_ITEM *),
757
journal_file_dict_items_backward_compar);
758
else
767
usec_t duration_ut = 0, max_duration_ut = 0;
768
usec_t progress_duration_ut = 0;
769
1156
- sampling_query_init(fqs, facets);
770
+ sampling_query_init(lqs, facets);
771
772
buffer_json_member_add_array(wb, "_journal_files");
773
for(size_t f = 0; f < files_used ;f++) {
774
const char *filename = dictionary_acquired_item_name(file_items[f]);
775
jf = dictionary_acquired_item_value(file_items[f]);
776
1163
- if(!jf_is_mine(jf, fqs))
777
+ if(!jf_is_mine(jf, lqs))
778
continue;
779
780
started_ut = ended_ut;
781
782
// do not even try to do the query if we expect it to pass the timeout
1169
- if(ended_ut + max_duration_ut * 3 >= *fqs->stop_monotonic_ut) {
783
+ if(ended_ut + max_duration_ut * 3 >= *lqs->stop_monotonic_ut) {
784
partial = true;
785
status = ND_SD_JOURNAL_TIMED_OUT;
786
break;
787
}
788
1175
- fqs->file_working++;
789
+ lqs->c.file_working++;
790
// fqs->cached_count = 0;
791
792
size_t fs_calls = fstat_thread_calls;
793
size_t fs_cached = fstat_thread_cached_responses;
1180
- size_t rows_useful = fqs->rows_useful;
1181
- size_t rows_read = fqs->rows_read;
1182
- size_t bytes_read = fqs->bytes_read;
1183
- size_t matches_setup_ut = fqs->matches_setup_ut;
794
+ size_t rows_useful = lqs->c.rows_useful;
795
+ size_t rows_read = lqs->c.rows_read;
796
+ size_t bytes_read = lqs->c.bytes_read;
797
+ size_t matches_setup_ut = lqs->c.matches_setup_ut;
798
1185
- sampling_file_init(fqs, jf);
799
+ sampling_file_init(lqs, jf);
800
1187
- ND_SD_JOURNAL_STATUS tmp_status = netdata_systemd_journal_query_one_file(filename, wb, facets, jf, fqs);
801
+ ND_SD_JOURNAL_STATUS tmp_status = netdata_systemd_journal_query_one_file(filename, wb, facets, jf, lqs);
802
803
// nd_log(NDLS_COLLECTORS, NDLP_INFO,
804
// "JOURNAL ESTIMATION FINAL: '%s' "
812
// , fqs->query_file.start_ut, fqs->query_file.stop_ut, fqs->query_file.stop_ut - fqs->query_file.start_ut
813
// );
814
1201
- rows_useful = fqs->rows_useful - rows_useful;
1202
- rows_read = fqs->rows_read - rows_read;
1203
- bytes_read = fqs->bytes_read - bytes_read;
1204
- matches_setup_ut = fqs->matches_setup_ut - matches_setup_ut;
815
+ rows_useful = lqs->c.rows_useful - rows_useful;
816
+ rows_read = lqs->c.rows_read - rows_read;
817
+ bytes_read = lqs->c.bytes_read - bytes_read;
818
+ matches_setup_ut = lqs->c.matches_setup_ut - matches_setup_ut;
819
fs_calls = fstat_thread_calls - fs_calls;
820
fs_cached = fstat_thread_cached_responses - fs_cached;
821
829
if(progress_duration_ut >= SYSTEMD_JOURNAL_PROGRESS_EVERY_UT) {
830
progress_duration_ut = 0;
831
netdata_mutex_lock(&stdout_mutex);
1218
- pluginsd_function_progress_to_stdout(fqs->transaction, f + 1, files_used);
832
+ pluginsd_function_progress_to_stdout(lqs->rq.transaction, f + 1, files_used);
833
netdata_mutex_unlock(&stdout_mutex);
834
}
835
855
buffer_json_member_add_uint64(wb, "fstat_query_calls", fs_calls);
856
buffer_json_member_add_uint64(wb, "fstat_query_cached_responses", fs_cached);
857
1244
- if(fqs->sampling) {
858
+ if(lqs->rq.sampling) {
859
buffer_json_member_add_object(wb, "_sampling");
860
{
1247
- buffer_json_member_add_uint64(wb, "sampled", fqs->samples_per_file.sampled);
1248
- buffer_json_member_add_uint64(wb, "unsampled", fqs->samples_per_file.unsampled);
1249
- buffer_json_member_add_uint64(wb, "estimated", fqs->samples_per_file.estimated);
861
+ buffer_json_member_add_uint64(wb, "sampled", lqs->c.samples_per_file.sampled);
862
+ buffer_json_member_add_uint64(wb, "unsampled", lqs->c.samples_per_file.unsampled);
863
+ buffer_json_member_add_uint64(wb, "estimated", lqs->c.samples_per_file.estimated);
864
}
865
buffer_json_object_close(wb); // _sampling
866
}
904
905
switch (status) {
906
case ND_SD_JOURNAL_OK:
1293
- if(fqs->if_modified_since && !fqs->rows_useful) {
1294
- buffer_flush(wb);
1295
- return HTTP_RESP_NOT_MODIFIED;
1296
- }
907
+ if(lqs->rq.if_modified_since && !lqs->c.rows_useful)
908
+ return rrd_call_function_error(wb, "no useful logs, not modified", HTTP_RESP_NOT_MODIFIED);
909
break;
910
911
case ND_SD_JOURNAL_TIMED_OUT:
913
break;
914
915
case ND_SD_JOURNAL_CANCELLED:
1304
- buffer_flush(wb);
1305
- return HTTP_RESP_CLIENT_CLOSED_REQUEST;
916
+ return rrd_call_function_error(wb, "client closed connection", HTTP_RESP_CLIENT_CLOSED_REQUEST);
917
918
case ND_SD_JOURNAL_NOT_MODIFIED:
1308
- buffer_flush(wb);
1309
- return HTTP_RESP_NOT_MODIFIED;
919
+ return rrd_call_function_error(wb, "not modified", HTTP_RESP_NOT_MODIFIED);
920
1311
- default:
921
case ND_SD_JOURNAL_FAILED_TO_OPEN:
922
+ return rrd_call_function_error(wb, "failed to open journal", HTTP_RESP_INTERNAL_SERVER_ERROR);
923
+
924
case ND_SD_JOURNAL_FAILED_TO_SEEK:
1314
- buffer_flush(wb);
1315
- return HTTP_RESP_INTERNAL_SERVER_ERROR;
925
+ return rrd_call_function_error(wb, "failed to seek in journal", HTTP_RESP_INTERNAL_SERVER_ERROR);
926
+
927
+ default:
928
+ return rrd_call_function_error(wb, "unknown status", HTTP_RESP_INTERNAL_SERVER_ERROR);
929
}
930
931
buffer_json_member_add_uint64(wb, "status", HTTP_RESP_OK);
933
buffer_json_member_add_string(wb, "type", "table");
934
935
// build a message for the query
1323
- if(!fqs->data_only) {
936
+ if(!lqs->rq.data_only) {
937
CLEAN_BUFFER *msg = buffer_create(0, NULL);
938
CLEAN_BUFFER *msg_description = buffer_create(0, NULL);
939
ND_LOG_FIELD_PRIORITY msg_priority = NDLP_INFO;
952
msg_priority = NDLP_WARNING;
953
}
954
1342
- if(fqs->samples.estimated || fqs->samples.unsampled) {
1343
- double percent = (double) (fqs->samples.sampled * 100.0 /
1344
- (fqs->samples.estimated + fqs->samples.unsampled + fqs->samples.sampled));
955
+ if(lqs->c.samples.estimated || lqs->c.samples.unsampled) {
956
+ double percent = (double) (lqs->c.samples.sampled * 100.0 /
957
+ (lqs->c.samples.estimated + lqs->c.samples.unsampled + lqs->c.samples.sampled));
958
buffer_sprintf(msg, "%.2f%% real data", percent);
959
buffer_sprintf(msg_description, "ACTUAL DATA: The filters counters reflect %0.2f%% of the data. ", percent);
960
msg_priority = MIN(msg_priority, NDLP_NOTICE);
961
}
962
1350
- if(fqs->samples.unsampled) {
1351
- double percent = (double) (fqs->samples.unsampled * 100.0 /
1352
- (fqs->samples.estimated + fqs->samples.unsampled + fqs->samples.sampled));
963
+ if(lqs->c.samples.unsampled) {
964
+ double percent = (double) (lqs->c.samples.unsampled * 100.0 /
965
+ (lqs->c.samples.estimated + lqs->c.samples.unsampled + lqs->c.samples.sampled));
966
buffer_sprintf(msg, ", %.2f%% unsampled", percent);
967
buffer_sprintf(msg_description
968
, "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. "
970
msg_priority = MIN(msg_priority, NDLP_NOTICE);
971
}
972
1360
- if(fqs->samples.estimated) {
1361
- double percent = (double) (fqs->samples.estimated * 100.0 /
1362
- (fqs->samples.estimated + fqs->samples.unsampled + fqs->samples.sampled));
973
+ if(lqs->c.samples.estimated) {
974
+ double percent = (double) (lqs->c.samples.estimated * 100.0 /
975
+ (lqs->c.samples.estimated + lqs->c.samples.unsampled + lqs->c.samples.sampled));
976
buffer_sprintf(msg, ", %.2f%% estimated", percent);
977
buffer_sprintf(msg_description
978
, "ESTIMATED DATA: The query selected a large amount of data, so to avoid delaying too much, the presented data are estimated by %0.2f%%. "
990
buffer_json_object_close(wb); // message
991
}
992
1380
- if(!fqs->data_only) {
993
+ if(!lqs->rq.data_only) {
994
buffer_json_member_add_time_t(wb, "update_every", 1);
995
buffer_json_member_add_string(wb, "help", SYSTEMD_JOURNAL_FUNCTION_DESCRIPTION);
996
}
997
1385
- if(!fqs->data_only || fqs->tail)
1386
- buffer_json_member_add_uint64(wb, "last_modified", fqs->last_modified);
998
+ if(!lqs->rq.data_only || lqs->rq.tail)
999
+ buffer_json_member_add_uint64(wb, "last_modified", lqs->last_modified);
1000
1001
facets_sort_and_reorder_keys(facets);
1002
facets_report(facets, wb, used_hashes_registry);
1003
1391
- buffer_json_member_add_time_t(wb, "expires", now_realtime_sec() + (fqs->data_only ? 3600 : 0));
1004
+ wb->expires = now_realtime_sec() + (lqs->rq.data_only ? 3600 : 0);
1005
+ buffer_json_member_add_time_t(wb, "expires", wb->expires);
1006
1007
buffer_json_member_add_object(wb, "_fstat_caching");
1008
{
1011
}
1012
buffer_json_object_close(wb); // _fstat_caching
1013
1400
- if(fqs->sampling) {
1014
+ if(lqs->rq.sampling) {
1015
buffer_json_member_add_object(wb, "_sampling");
1016
{
1403
- buffer_json_member_add_uint64(wb, "sampled", fqs->samples.sampled);
1404
- buffer_json_member_add_uint64(wb, "unsampled", fqs->samples.unsampled);
1405
- buffer_json_member_add_uint64(wb, "estimated", fqs->samples.estimated);
1017
+ buffer_json_member_add_uint64(wb, "sampled", lqs->c.samples.sampled);
1018
+ buffer_json_member_add_uint64(wb, "unsampled", lqs->c.samples.unsampled);
1019
+ buffer_json_member_add_uint64(wb, "estimated", lqs->c.samples.estimated);
1020
}
1021
buffer_json_object_close(wb); // _sampling
1022
}
1023
1410
- buffer_json_finalize(wb);
1411
-
1412
- return HTTP_RESP_OK;
1024
+ wb->content_type = CT_APPLICATION_JSON;
1025
+ wb->response_code = HTTP_RESP_OK;
1026
+ return wb->response_code;
1027
}
1028
1415
-static void netdata_systemd_journal_function_help(const char *transaction) {
1416
- CLEAN_BUFFER *wb = buffer_create(0, NULL);
1417
- buffer_sprintf(wb,
1418
- "%s / %s\n"
1419
- "\n"
1420
- "%s\n"
1421
- "\n"
1422
- "The following parameters are supported:\n"
1423
- "\n"
1424
- " "JOURNAL_PARAMETER_HELP"\n"
1425
- " Shows this help message.\n"
1426
- "\n"
1427
- " "JOURNAL_PARAMETER_INFO"\n"
1428
- " Request initial configuration information about the plugin.\n"
1429
- " The key entity returned is the required_params array, which includes\n"
1430
- " all the available systemd journal sources.\n"
1431
- " When `"JOURNAL_PARAMETER_INFO"` is requested, all other parameters are ignored.\n"
1432
- "\n"
1433
- " "JOURNAL_PARAMETER_DATA_ONLY":true or "JOURNAL_PARAMETER_DATA_ONLY":false\n"
1434
- " Quickly respond with data requested, without generating a\n"
1435
- " `histogram`, `facets` counters and `items`.\n"
1436
- "\n"
1437
- " "JOURNAL_PARAMETER_DELTA":true or "JOURNAL_PARAMETER_DELTA":false\n"
1438
- " When doing data only queries, include deltas for histogram, facets and items.\n"
1439
- "\n"
1440
- " "JOURNAL_PARAMETER_TAIL":true or "JOURNAL_PARAMETER_TAIL":false\n"
1441
- " When doing data only queries, respond with the newest messages,\n"
1442
- " and up to the anchor, but calculate deltas (if requested) for\n"
1443
- " the duration [anchor - before].\n"
1444
- "\n"
1445
- " "JOURNAL_PARAMETER_SLICE":true or "JOURNAL_PARAMETER_SLICE":false\n"
1446
- " When it is turned on, the plugin is executing filtering via libsystemd,\n"
1447
- " utilizing all the available indexes of the journal files.\n"
1448
- " When it is off, only the time constraint is handled by libsystemd and\n"
1449
- " all filtering is done by the plugin.\n"
1450
- " The default is: %s\n"
1451
- "\n"
1452
- " "JOURNAL_PARAMETER_SOURCE":SOURCE\n"
1453
- " Query only the specified journal sources.\n"
1454
- " Do an `"JOURNAL_PARAMETER_INFO"` query to find the sources.\n"
1455
- "\n"
1456
- " "JOURNAL_PARAMETER_BEFORE":TIMESTAMP_IN_SECONDS\n"
1457
- " Absolute or relative (to now) timestamp in seconds, to start the query.\n"
1458
- " The query is always executed from the most recent to the oldest log entry.\n"
1459
- " If not given the default is: now.\n"
1460
- "\n"
1461
- " "JOURNAL_PARAMETER_AFTER":TIMESTAMP_IN_SECONDS\n"
1462
- " Absolute or relative (to `before`) timestamp in seconds, to end the query.\n"
1463
- " If not given, the default is %d.\n"
1464
- "\n"
1465
- " "JOURNAL_PARAMETER_LAST":ITEMS\n"
1466
- " The number of items to return.\n"
1467
- " The default is %d.\n"
1468
- "\n"
1469
- " "JOURNAL_PARAMETER_SAMPLING":ITEMS\n"
1470
- " The number of log entries to sample to estimate facets counters and histogram.\n"
1471
- " The default is %d.\n"
1472
- "\n"
1473
- " "JOURNAL_PARAMETER_ANCHOR":TIMESTAMP_IN_MICROSECONDS\n"
1474
- " Return items relative to this timestamp.\n"
1475
- " The exact items to be returned depend on the query `"JOURNAL_PARAMETER_DIRECTION"`.\n"
1476
- "\n"
1477
- " "JOURNAL_PARAMETER_DIRECTION":forward or "JOURNAL_PARAMETER_DIRECTION":backward\n"
1478
- " When set to `backward` (default) the items returned are the newest before the\n"
1479
- " `"JOURNAL_PARAMETER_ANCHOR"`, (or `"JOURNAL_PARAMETER_BEFORE"` if `"JOURNAL_PARAMETER_ANCHOR"` is not set)\n"
1480
- " When set to `forward` the items returned are the oldest after the\n"
1481
- " `"JOURNAL_PARAMETER_ANCHOR"`, (or `"JOURNAL_PARAMETER_AFTER"` if `"JOURNAL_PARAMETER_ANCHOR"` is not set)\n"
1482
- " The default is: %s\n"
1483
- "\n"
1484
- " "JOURNAL_PARAMETER_QUERY":SIMPLE_PATTERN\n"
1485
- " Do a full text search to find the log entries matching the pattern given.\n"
1486
- " The plugin is searching for matches on all fields of the database.\n"
1487
- "\n"
1488
- " "JOURNAL_PARAMETER_IF_MODIFIED_SINCE":TIMESTAMP_IN_MICROSECONDS\n"
1489
- " Each successful response, includes a `last_modified` field.\n"
1490
- " By providing the timestamp to the `"JOURNAL_PARAMETER_IF_MODIFIED_SINCE"` parameter,\n"
1491
- " the plugin will return 200 with a successful response, or 304 if the source has not\n"
1492
- " been modified since that timestamp.\n"
1493
- "\n"
1494
- " "JOURNAL_PARAMETER_HISTOGRAM":facet_id\n"
1495
- " Use the given `facet_id` for the histogram.\n"
1496
- " This parameter is ignored in `"JOURNAL_PARAMETER_DATA_ONLY"` mode.\n"
1497
- "\n"
1498
- " "JOURNAL_PARAMETER_FACETS":facet_id1,facet_id2,facet_id3,...\n"
1499
- " Add the given facets to the list of fields for which analysis is required.\n"
1500
- " The plugin will offer both a histogram and facet value counters for its values.\n"
1501
- " This parameter is ignored in `"JOURNAL_PARAMETER_DATA_ONLY"` mode.\n"
1502
- "\n"
1503
- " facet_id:value_id1,value_id2,value_id3,...\n"
1504
- " Apply filters to the query, based on the facet IDs returned.\n"
1505
- " Each `facet_id` can be given once, but multiple `facet_ids` can be given.\n"
1506
- "\n"
1507
- , program_name
1508
- , SYSTEMD_JOURNAL_FUNCTION_NAME
1509
- , SYSTEMD_JOURNAL_FUNCTION_DESCRIPTION
1510
- , JOURNAL_DEFAULT_SLICE_MODE ? "true" : "false" // slice
1511
- , -SYSTEMD_JOURNAL_DEFAULT_QUERY_DURATION
1512
- , SYSTEMD_JOURNAL_DEFAULT_ITEMS_PER_QUERY
1513
- , SYSTEMD_JOURNAL_DEFAULT_ITEMS_SAMPLING
1514
- , JOURNAL_DEFAULT_DIRECTION == FACETS_ANCHOR_DIRECTION_BACKWARD ? "backward" : "forward"
1515
- );
1516
-
1517
- netdata_mutex_lock(&stdout_mutex);
1518
- pluginsd_function_result_to_stdout(transaction, HTTP_RESP_OK, "text/plain", now_realtime_sec() + 3600, wb);
1519
- netdata_mutex_unlock(&stdout_mutex);
1520
-}
1521
-
1522
-typedef struct {
1523
- FACET_KEY_OPTIONS default_facet;
1524
- bool info;
1525
- bool data_only;
1526
- bool slice;
1527
- bool delta;
1528
- bool tail;
1529
- time_t after_s;
1530
- time_t before_s;
1531
- usec_t anchor;
1532
- usec_t if_modified_since;
1533
- size_t last;
1534
- FACETS_ANCHOR_DIRECTION direction;
1535
- const char *query;
1536
- const char *chart;
1537
- SIMPLE_PATTERN *sources;
1538
- SD_JOURNAL_FILE_SOURCE_TYPE source_type;
1539
- size_t filters;
1540
- size_t sampling;
1541
-} JOURNAL_QUERY;
1542
-
1543
-static SD_JOURNAL_FILE_SOURCE_TYPE get_internal_source_type(const char *value) {
1544
- if(strcmp(value, SDJF_SOURCE_ALL_NAME) == 0)
1545
- return SDJF_ALL;
1546
- else if(strcmp(value, SDJF_SOURCE_LOCAL_NAME) == 0)
1547
- return SDJF_LOCAL_ALL;
1548
- else if(strcmp(value, SDJF_SOURCE_REMOTES_NAME) == 0)
1549
- return SDJF_REMOTE_ALL;
1550
- else if(strcmp(value, SDJF_SOURCE_NAMESPACES_NAME) == 0)
1551
- return SDJF_LOCAL_NAMESPACE;
1552
- else if(strcmp(value, SDJF_SOURCE_LOCAL_SYSTEM_NAME) == 0)
1553
- return SDJF_LOCAL_SYSTEM;
1554
- else if(strcmp(value, SDJF_SOURCE_LOCAL_USERS_NAME) == 0)
1555
- return SDJF_LOCAL_USER;
1556
- else if(strcmp(value, SDJF_SOURCE_LOCAL_OTHER_NAME) == 0)
1557
- return SDJF_LOCAL_OTHER;
1558
-
1559
- return SDJF_NONE;
1560
-}
1561
-
1562
-static FACETS_ANCHOR_DIRECTION get_direction(const char *value) {
1563
- return strcasecmp(value, "forward") == 0 ? FACETS_ANCHOR_DIRECTION_FORWARD : FACETS_ANCHOR_DIRECTION_BACKWARD;
1564
-}
1565
-
1566
-struct post_query_data {
1567
- const char *transaction;
1568
- FACETS *facets;
1569
- JOURNAL_QUERY *q;
1570
- BUFFER *wb;
1571
-};
1572
-
1573
-static bool parse_json_payload(json_object *jobj, const char *path, void *data, BUFFER *error) {
1574
- struct post_query_data *qd = data;
1575
- JOURNAL_QUERY *q = qd->q;
1576
- BUFFER *wb = qd->wb;
1577
- FACETS *facets = qd->facets;
1578
- // const char *transaction = qd->transaction;
1579
-
1580
- buffer_flush(error);
1581
-
1582
- JSONC_PARSE_BOOL_OR_ERROR_AND_RETURN(jobj, path, JOURNAL_PARAMETER_INFO, q->info, error, false);
1583
- JSONC_PARSE_BOOL_OR_ERROR_AND_RETURN(jobj, path, JOURNAL_PARAMETER_DELTA, q->delta, error, false);
1584
- JSONC_PARSE_BOOL_OR_ERROR_AND_RETURN(jobj, path, JOURNAL_PARAMETER_TAIL, q->tail, error, false);
1585
- JSONC_PARSE_BOOL_OR_ERROR_AND_RETURN(jobj, path, JOURNAL_PARAMETER_SLICE, q->slice, error, false);
1586
- JSONC_PARSE_BOOL_OR_ERROR_AND_RETURN(jobj, path, JOURNAL_PARAMETER_DATA_ONLY, q->data_only, error, false);
1587
- JSONC_PARSE_UINT64_OR_ERROR_AND_RETURN(jobj, path, JOURNAL_PARAMETER_SAMPLING, q->sampling, error, false);
1588
- JSONC_PARSE_INT64_OR_ERROR_AND_RETURN(jobj, path, JOURNAL_PARAMETER_AFTER, q->after_s, error, false);
1589
- JSONC_PARSE_INT64_OR_ERROR_AND_RETURN(jobj, path, JOURNAL_PARAMETER_BEFORE, q->before_s, error, false);
1590
- JSONC_PARSE_UINT64_OR_ERROR_AND_RETURN(jobj, path, JOURNAL_PARAMETER_IF_MODIFIED_SINCE, q->if_modified_since, error, false);
1591
- JSONC_PARSE_UINT64_OR_ERROR_AND_RETURN(jobj, path, JOURNAL_PARAMETER_ANCHOR, q->anchor, error, false);
1592
- JSONC_PARSE_UINT64_OR_ERROR_AND_RETURN(jobj, path, JOURNAL_PARAMETER_LAST, q->last, error, false);
1593
- JSONC_PARSE_TXT2ENUM_OR_ERROR_AND_RETURN(jobj, path, JOURNAL_PARAMETER_DIRECTION, get_direction, q->direction, error, false);
1594
- JSONC_PARSE_TXT2STRDUPZ_OR_ERROR_AND_RETURN(jobj, path, JOURNAL_PARAMETER_QUERY, q->query, error, false);
1595
- JSONC_PARSE_TXT2STRDUPZ_OR_ERROR_AND_RETURN(jobj, path, JOURNAL_PARAMETER_HISTOGRAM, q->chart, error, false);
1596
-
1597
- json_object *sources;
1598
- if (json_object_object_get_ex(jobj, JOURNAL_PARAMETER_SOURCE, &sources)) {
1599
- if (json_object_get_type(sources) != json_type_array) {
1600
- buffer_sprintf(error, "member '%s' is not an array", JOURNAL_PARAMETER_SOURCE);
1601
- return false;
1602
- }
1603
-
1604
- buffer_json_member_add_array(wb, JOURNAL_PARAMETER_SOURCE);
1605
-
1606
- CLEAN_BUFFER *sources_list = buffer_create(0, NULL);
1607
-
1608
- q->source_type = SDJF_NONE;
1609
-
1610
- size_t sources_len = json_object_array_length(sources);
1611
- for (size_t i = 0; i < sources_len; i++) {
1612
- json_object *src = json_object_array_get_idx(sources, i);
1613
-
1614
- if (json_object_get_type(src) != json_type_string) {
1615
- buffer_sprintf(error, "sources array item %zu is not a string", i);
1616
- return false;
1617
- }
1618
-
1619
- const char *value = json_object_get_string(src);
1620
- buffer_json_add_array_item_string(wb, value);
1621
-
1622
- SD_JOURNAL_FILE_SOURCE_TYPE t = get_internal_source_type(value);
1623
- if(t != SDJF_NONE) {
1624
- q->source_type |= t;
1625
- value = NULL;
1626
- }
1627
- else {
1628
- // else, match the source, whatever it is
1629
- if(buffer_strlen(sources_list))
1630
- buffer_putc(sources_list, '|');
1631
-
1632
- buffer_strcat(sources_list, value);
1633
- }
1634
- }
1635
-
1636
- if(buffer_strlen(sources_list)) {
1637
- simple_pattern_free(q->sources);
1638
- q->sources = simple_pattern_create(buffer_tostring(sources_list), "|", SIMPLE_PATTERN_EXACT, false);
1639
- }
1640
-
1641
- buffer_json_array_close(wb); // source
1642
- }
1643
-
1644
- json_object *fcts;
1645
- if (json_object_object_get_ex(jobj, JOURNAL_PARAMETER_FACETS, &fcts)) {
1646
- if (json_object_get_type(sources) != json_type_array) {
1647
- buffer_sprintf(error, "member '%s' is not an array", JOURNAL_PARAMETER_FACETS);
1648
- return false;
1649
- }
1650
-
1651
- q->default_facet = FACET_KEY_OPTION_NONE;
1652
- facets_reset_and_disable_all_facets(facets);
1653
-
1654
- buffer_json_member_add_array(wb, JOURNAL_PARAMETER_FACETS);
1655
-
1656
- size_t facets_len = json_object_array_length(fcts);
1657
- for (size_t i = 0; i < facets_len; i++) {
1658
- json_object *fct = json_object_array_get_idx(fcts, i);
1659
-
1660
- if (json_object_get_type(fct) != json_type_string) {
1661
- buffer_sprintf(error, "facets array item %zu is not a string", i);
1662
- return false;
1663
- }
1664
-
1665
- const char *value = json_object_get_string(fct);
1666
- facets_register_facet(facets, value, FACET_KEY_OPTION_FACET|FACET_KEY_OPTION_FTS|FACET_KEY_OPTION_REORDER);
1667
- buffer_json_add_array_item_string(wb, value);
1668
- }
1669
-
1670
- buffer_json_array_close(wb); // facets
1671
- }
1672
-
1673
- json_object *selections;
1674
- if (json_object_object_get_ex(jobj, "selections", &selections)) {
1675
- if (json_object_get_type(selections) != json_type_object) {
1676
- buffer_sprintf(error, "member 'selections' is not an object");
1677
- return false;
1678
- }
1679
-
1680
- buffer_json_member_add_object(wb, "selections");
1681
-
1682
- json_object_object_foreach(selections, key, val) {
1683
- if (json_object_get_type(val) != json_type_array) {
1684
- buffer_sprintf(error, "selection '%s' is not an array", key);
1685
- return false;
1686
- }
1687
-
1688
- buffer_json_member_add_array(wb, key);
1689
-
1690
- size_t values_len = json_object_array_length(val);
1691
- for (size_t i = 0; i < values_len; i++) {
1692
- json_object *value_obj = json_object_array_get_idx(val, i);
1693
-
1694
- if (json_object_get_type(value_obj) != json_type_string) {
1695
- buffer_sprintf(error, "selection '%s' array item %zu is not a string", key, i);
1696
- return false;
1697
- }
1698
-
1699
- const char *value = json_object_get_string(value_obj);
1700
-
1701
- // Call facets_register_facet_id_filter for each value
1702
- facets_register_facet_filter(
1703
- facets, key, value, FACET_KEY_OPTION_FACET | FACET_KEY_OPTION_FTS | FACET_KEY_OPTION_REORDER);
1704
-
1705
- buffer_json_add_array_item_string(wb, value);
1706
- q->filters++;
1707
- }
1708
-
1709
- buffer_json_array_close(wb); // key
1710
- }
1711
-
1712
- buffer_json_object_close(wb); // selections
1713
- }
1714
-
1715
- return true;
1716
-}
1717
-
1718
-static bool parse_post_params(FACETS *facets, JOURNAL_QUERY *q, BUFFER *wb, BUFFER *payload, const char *transaction) {
1719
- struct post_query_data qd = {
1720
- .transaction = transaction,
1721
- .facets = facets,
1722
- .q = q,
1723
- .wb = wb,
1724
- };
1725
-
1726
- int code;
1727
- CLEAN_JSON_OBJECT *jobj = json_parse_function_payload_or_error(wb, payload, &code, parse_json_payload, &qd);
1728
- if(!jobj || code != HTTP_RESP_OK) {
1729
- netdata_mutex_lock(&stdout_mutex);
1730
- pluginsd_function_result_to_stdout(transaction, code, "application/json", now_realtime_sec() + 1, wb);
1731
- netdata_mutex_unlock(&stdout_mutex);
1732
- return false;
1733
- }
1734
-
1735
- return true;
1736
-}
1737
-
1738
-static bool parse_get_params(FACETS *facets, JOURNAL_QUERY *q, BUFFER *wb, char *function, const char *transaction) {
1739
- buffer_json_member_add_object(wb, "_request");
1740
-
1741
- char *words[SYSTEMD_JOURNAL_MAX_PARAMS] = { NULL };
1742
- size_t num_words = quoted_strings_splitter_pluginsd(function, words, SYSTEMD_JOURNAL_MAX_PARAMS);
1743
- for(int i = 1; i < SYSTEMD_JOURNAL_MAX_PARAMS ;i++) {
1744
- char *keyword = get_word(words, num_words, i);
1745
- if(!keyword) break;
1746
-
1747
- if(strcmp(keyword, JOURNAL_PARAMETER_HELP) == 0) {
1748
- netdata_systemd_journal_function_help(transaction);
1749
- return false;
1750
- }
1751
- else if(strcmp(keyword, JOURNAL_PARAMETER_INFO) == 0) {
1752
- q->info = true;
1753
- }
1754
- else if(strncmp(keyword, JOURNAL_PARAMETER_DELTA ":", sizeof(JOURNAL_PARAMETER_DELTA ":") - 1) == 0) {
1755
- char *v = &keyword[sizeof(JOURNAL_PARAMETER_DELTA ":") - 1];
1756
-
1757
- if(strcmp(v, "false") == 0 || strcmp(v, "no") == 0 || strcmp(v, "0") == 0)
1758
- q->delta = false;
1759
- else
1760
- q->delta = true;
1761
- }
1762
- else if(strncmp(keyword, JOURNAL_PARAMETER_TAIL ":", sizeof(JOURNAL_PARAMETER_TAIL ":") - 1) == 0) {
1763
- char *v = &keyword[sizeof(JOURNAL_PARAMETER_TAIL ":") - 1];
1764
-
1765
- if(strcmp(v, "false") == 0 || strcmp(v, "no") == 0 || strcmp(v, "0") == 0)
1766
- q->tail = false;
1767
- else
1768
- q->tail = true;
1769
- }
1770
- else if(strncmp(keyword, JOURNAL_PARAMETER_SAMPLING ":", sizeof(JOURNAL_PARAMETER_SAMPLING ":") - 1) == 0) {
1771
- q->sampling = str2ul(&keyword[sizeof(JOURNAL_PARAMETER_SAMPLING ":") - 1]);
1772
- }
1773
- else if(strncmp(keyword, JOURNAL_PARAMETER_DATA_ONLY ":", sizeof(JOURNAL_PARAMETER_DATA_ONLY ":") - 1) == 0) {
1774
- char *v = &keyword[sizeof(JOURNAL_PARAMETER_DATA_ONLY ":") - 1];
1775
-
1776
- if(strcmp(v, "false") == 0 || strcmp(v, "no") == 0 || strcmp(v, "0") == 0)
1777
- q->data_only = false;
1778
- else
1779
- q->data_only = true;
1780
- }
1781
- else if(strncmp(keyword, JOURNAL_PARAMETER_SLICE ":", sizeof(JOURNAL_PARAMETER_SLICE ":") - 1) == 0) {
1782
- char *v = &keyword[sizeof(JOURNAL_PARAMETER_SLICE ":") - 1];
1783
-
1784
- if(strcmp(v, "false") == 0 || strcmp(v, "no") == 0 || strcmp(v, "0") == 0)
1785
- q->slice = false;
1786
- else
1787
- q->slice = true;
1788
- }
1789
- else if(strncmp(keyword, JOURNAL_PARAMETER_SOURCE ":", sizeof(JOURNAL_PARAMETER_SOURCE ":") - 1) == 0) {
1790
- const char *value = &keyword[sizeof(JOURNAL_PARAMETER_SOURCE ":") - 1];
1791
-
1792
- buffer_json_member_add_array(wb, JOURNAL_PARAMETER_SOURCE);
1793
-
1794
- CLEAN_BUFFER *sources_list = buffer_create(0, NULL);
1795
-
1796
- q->source_type = SDJF_NONE;
1797
- while(value) {
1798
- char *sep = strchr(value, ',');
1799
- if(sep)
1800
- *sep++ = '\0';
1801
-
1802
- buffer_json_add_array_item_string(wb, value);
1803
-
1804
- SD_JOURNAL_FILE_SOURCE_TYPE t = get_internal_source_type(value);
1805
- if(t != SDJF_NONE) {
1806
- q->source_type |= t;
1807
- value = NULL;
1808
- }
1809
- else {
1810
- // else, match the source, whatever it is
1811
- if(buffer_strlen(sources_list))
1812
- buffer_putc(sources_list, '|');
1813
-
1814
- buffer_strcat(sources_list, value);
1815
- }
1816
-
1817
- value = sep;
1818
- }
1819
-
1820
- if(buffer_strlen(sources_list)) {
1821
- simple_pattern_free(q->sources);
1822
- q->sources = simple_pattern_create(buffer_tostring(sources_list), "|", SIMPLE_PATTERN_EXACT, false);
1823
- }
1824
-
1825
- buffer_json_array_close(wb); // source
1826
- }
1827
- else if(strncmp(keyword, JOURNAL_PARAMETER_AFTER ":", sizeof(JOURNAL_PARAMETER_AFTER ":") - 1) == 0) {
1828
- q->after_s = str2l(&keyword[sizeof(JOURNAL_PARAMETER_AFTER ":") - 1]);
1829
- }
1830
- else if(strncmp(keyword, JOURNAL_PARAMETER_BEFORE ":", sizeof(JOURNAL_PARAMETER_BEFORE ":") - 1) == 0) {
1831
- q->before_s = str2l(&keyword[sizeof(JOURNAL_PARAMETER_BEFORE ":") - 1]);
1832
- }
1833
- else if(strncmp(keyword, JOURNAL_PARAMETER_IF_MODIFIED_SINCE ":", sizeof(JOURNAL_PARAMETER_IF_MODIFIED_SINCE ":") - 1) == 0) {
1834
- q->if_modified_since = str2ull(&keyword[sizeof(JOURNAL_PARAMETER_IF_MODIFIED_SINCE ":") - 1], NULL);
1835
- }
1836
- else if(strncmp(keyword, JOURNAL_PARAMETER_ANCHOR ":", sizeof(JOURNAL_PARAMETER_ANCHOR ":") - 1) == 0) {
1837
- q->anchor = str2ull(&keyword[sizeof(JOURNAL_PARAMETER_ANCHOR ":") - 1], NULL);
1838
- }
1839
- else if(strncmp(keyword, JOURNAL_PARAMETER_DIRECTION ":", sizeof(JOURNAL_PARAMETER_DIRECTION ":") - 1) == 0) {
1840
- q->direction = get_direction(&keyword[sizeof(JOURNAL_PARAMETER_DIRECTION ":") - 1]);
1841
- }
1842
- else if(strncmp(keyword, JOURNAL_PARAMETER_LAST ":", sizeof(JOURNAL_PARAMETER_LAST ":") - 1) == 0) {
1843
- q->last = str2ul(&keyword[sizeof(JOURNAL_PARAMETER_LAST ":") - 1]);
1844
- }
1845
- else if(strncmp(keyword, JOURNAL_PARAMETER_QUERY ":", sizeof(JOURNAL_PARAMETER_QUERY ":") - 1) == 0) {
1846
- freez((void *)q->query);
1847
- q->query= strdupz(&keyword[sizeof(JOURNAL_PARAMETER_QUERY ":") - 1]);
1848
- }
1849
- else if(strncmp(keyword, JOURNAL_PARAMETER_HISTOGRAM ":", sizeof(JOURNAL_PARAMETER_HISTOGRAM ":") - 1) == 0) {
1850
- freez((void *)q->chart);
1851
- q->chart = strdupz(&keyword[sizeof(JOURNAL_PARAMETER_HISTOGRAM ":") - 1]);
1852
- }
1853
- else if(strncmp(keyword, JOURNAL_PARAMETER_FACETS ":", sizeof(JOURNAL_PARAMETER_FACETS ":") - 1) == 0) {
1854
- q->default_facet = FACET_KEY_OPTION_NONE;
1855
- facets_reset_and_disable_all_facets(facets);
1856
-
1857
- char *value = &keyword[sizeof(JOURNAL_PARAMETER_FACETS ":") - 1];
1858
- if(*value) {
1859
- buffer_json_member_add_array(wb, JOURNAL_PARAMETER_FACETS);
1860
-
1861
- while(value) {
1862
- char *sep = strchr(value, ',');
1863
- if(sep)
1864
- *sep++ = '\0';
1865
-
1866
- facets_register_facet_id(facets, value, FACET_KEY_OPTION_FACET|FACET_KEY_OPTION_FTS|FACET_KEY_OPTION_REORDER);
1867
- buffer_json_add_array_item_string(wb, value);
1868
-
1869
- value = sep;
1870
- }
1871
-
1872
- buffer_json_array_close(wb); // JOURNAL_PARAMETER_FACETS
1873
- }
1874
- }
1875
- else {
1876
- char *value = strchr(keyword, ':');
1877
- if(value) {
1878
- *value++ = '\0';
1879
-
1880
- buffer_json_member_add_array(wb, keyword);
1881
-
1882
- while(value) {
1883
- char *sep = strchr(value, ',');
1884
- if(sep)
1885
- *sep++ = '\0';
1886
-
1887
- facets_register_facet_filter_id(
1888
- facets,
1889
- keyword,
1890
- value,
1891
- FACET_KEY_OPTION_FACET | FACET_KEY_OPTION_FTS | FACET_KEY_OPTION_REORDER);
1892
-
1893
- buffer_json_add_array_item_string(wb, value);
1894
- q->filters++;
1895
-
1896
- value = sep;
1897
- }
1898
-
1899
- buffer_json_array_close(wb); // keyword
1900
- }
1901
- }
1902
- }
1903
-
1904
- return true;
1905
-}
1906
-
1907
-void function_systemd_journal(const char *transaction, char *function, usec_t *stop_monotonic_ut, bool *cancelled,
1908
- BUFFER *payload, HTTP_ACCESS access __maybe_unused,
1909
- const char *source __maybe_unused, void *data __maybe_unused) {
1910
- fstat_thread_calls = 0;
1911
- fstat_thread_cached_responses = 0;
1912
-
1913
- CLEAN_BUFFER *wb = buffer_create(0, NULL);
1914
- buffer_flush(wb);
1915
- buffer_json_initialize(wb, "\"", "\"", 0, true, BUFFER_JSON_OPTIONS_MINIFY);
1916
-
1917
- FUNCTION_QUERY_STATUS tmp_fqs = {
1918
- .cancelled = cancelled,
1919
- .stop_monotonic_ut = stop_monotonic_ut,
1920
- };
1921
- FUNCTION_QUERY_STATUS *fqs = NULL;
1922
-
1923
- FACETS *facets = facets_create(50, FACETS_OPTION_ALL_KEYS_FTS,
1924
- SYSTEMD_ALWAYS_VISIBLE_KEYS,
1925
- SYSTEMD_KEYS_INCLUDED_IN_FACETS,
1926
- SYSTEMD_KEYS_EXCLUDED_FROM_FACETS);
1927
-
1928
- facets_accepted_param(facets, JOURNAL_PARAMETER_INFO);
1929
- facets_accepted_param(facets, JOURNAL_PARAMETER_SOURCE);
1930
- facets_accepted_param(facets, JOURNAL_PARAMETER_AFTER);
1931
- facets_accepted_param(facets, JOURNAL_PARAMETER_BEFORE);
1932
- facets_accepted_param(facets, JOURNAL_PARAMETER_ANCHOR);
1933
- facets_accepted_param(facets, JOURNAL_PARAMETER_DIRECTION);
1934
- facets_accepted_param(facets, JOURNAL_PARAMETER_LAST);
1935
- facets_accepted_param(facets, JOURNAL_PARAMETER_QUERY);
1936
- facets_accepted_param(facets, JOURNAL_PARAMETER_FACETS);
1937
- facets_accepted_param(facets, JOURNAL_PARAMETER_HISTOGRAM);
1938
- facets_accepted_param(facets, JOURNAL_PARAMETER_IF_MODIFIED_SINCE);
1939
- facets_accepted_param(facets, JOURNAL_PARAMETER_DATA_ONLY);
1940
- facets_accepted_param(facets, JOURNAL_PARAMETER_DELTA);
1941
- facets_accepted_param(facets, JOURNAL_PARAMETER_TAIL);
1942
- facets_accepted_param(facets, JOURNAL_PARAMETER_SAMPLING);
1943
-
1944
-#ifdef HAVE_SD_JOURNAL_RESTART_FIELDS
1945
- facets_accepted_param(facets, JOURNAL_PARAMETER_SLICE);
1946
-#endif // HAVE_SD_JOURNAL_RESTART_FIELDS
1947
-
1948
- // ------------------------------------------------------------------------
1949
- // parse the parameters
1950
-
1951
- JOURNAL_QUERY q = {
1952
- .default_facet = FACET_KEY_OPTION_FACET,
1953
- .info = false,
1954
- .data_only = false,
1955
- .slice = JOURNAL_DEFAULT_SLICE_MODE,
1956
- .delta = false,
1957
- .tail = false,
1958
- .after_s = 0,
1959
- .before_s = 0,
1960
- .anchor = 0,
1961
- .if_modified_since = 0,
1962
- .last = 0,
1963
- .direction = JOURNAL_DEFAULT_DIRECTION,
1964
- .query = NULL,
1965
- .chart = NULL,
1966
- .sources = NULL,
1967
- .source_type = SDJF_ALL,
1968
- .filters = 0,
1969
- .sampling = SYSTEMD_JOURNAL_DEFAULT_ITEMS_SAMPLING,
1970
- };
1971
-
1972
- if( (payload && !parse_post_params(facets, &q, wb, payload, transaction)) ||
1973
- (!payload && !parse_get_params(facets, &q, wb, function, transaction)) )
1974
- goto cleanup;
1029
+static void systemd_journal_register_transformations(LOGS_QUERY_STATUS *lqs) {
1030
+ FACETS *facets = lqs->facets;
1031
+ LOGS_QUERY_REQUEST *rq = &lqs->rq;
1032
1033
// ----------------------------------------------------------------------------------------------------------------
1034
// register the fields in the order you want them on the dashboard
1036
facets_register_row_severity(facets, syslog_priority_to_facet_severity, NULL);
1037
1038
facets_register_key_name(
1982
- facets, "_HOSTNAME",
1983
- q.default_facet | FACET_KEY_OPTION_VISIBLE);
1039
+ facets, "_HOSTNAME", rq->default_facet | FACET_KEY_OPTION_VISIBLE);
1040
1041
facets_register_dynamic_key_name(
1042
facets, JOURNAL_KEY_ND_JOURNAL_PROCESS,
1056
1057
facets_register_key_name_transformation(
1058
facets, "PRIORITY",
2003
- q.default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW |
1059
+ rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW |
1060
FACET_KEY_OPTION_EXPANDED_FILTER,
1061
netdata_systemd_journal_transform_priority, NULL);
1062
1063
facets_register_key_name_transformation(
1064
facets, "SYSLOG_FACILITY",
2009
- q.default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW |
1065
+ rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW |
1066
FACET_KEY_OPTION_EXPANDED_FILTER,
1067
netdata_systemd_journal_transform_syslog_facility, NULL);
1068
1069
facets_register_key_name_transformation(
1070
facets, "ERRNO",
2015
- q.default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1071
+ rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1072
netdata_systemd_journal_transform_errno, NULL);
1073
1074
facets_register_key_name(
1076
FACET_KEY_OPTION_NEVER_FACET);
1077
1078
facets_register_key_name(
2023
- facets, "SYSLOG_IDENTIFIER",
2024
- q.default_facet);
1079
+ facets, "SYSLOG_IDENTIFIER", rq->default_facet);
1080
1081
facets_register_key_name(
2027
- facets, "UNIT",
2028
- q.default_facet);
1082
+ facets, "UNIT", rq->default_facet);
1083
1084
facets_register_key_name(
2031
- facets, "USER_UNIT",
2032
- q.default_facet);
1085
+ facets, "USER_UNIT", rq->default_facet);
1086
1087
facets_register_key_name_transformation(
1088
facets, "MESSAGE_ID",
2036
- q.default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW |
1089
+ rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW |
1090
FACET_KEY_OPTION_EXPANDED_FILTER,
1091
netdata_systemd_journal_transform_message_id, NULL);
1092
1093
facets_register_key_name_transformation(
1094
facets, "_BOOT_ID",
2042
- q.default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1095
+ rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1096
netdata_systemd_journal_transform_boot_id, NULL);
1097
1098
facets_register_key_name_transformation(
1099
facets, "_SYSTEMD_OWNER_UID",
2047
- q.default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1100
+ rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1101
netdata_systemd_journal_transform_uid, NULL);
1102
1103
facets_register_key_name_transformation(
1104
facets, "_UID",
2052
- q.default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1105
+ rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1106
netdata_systemd_journal_transform_uid, NULL);
1107
1108
facets_register_key_name_transformation(
1109
facets, "OBJECT_SYSTEMD_OWNER_UID",
2057
- q.default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1110
+ rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1111
netdata_systemd_journal_transform_uid, NULL);
1112
1113
facets_register_key_name_transformation(
1114
facets, "OBJECT_UID",
2062
- q.default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1115
+ rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1116
netdata_systemd_journal_transform_uid, NULL);
1117
1118
facets_register_key_name_transformation(
1119
facets, "_GID",
2067
- q.default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1120
+ rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1121
netdata_systemd_journal_transform_gid, NULL);
1122
1123
facets_register_key_name_transformation(
1124
facets, "OBJECT_GID",
2072
- q.default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1125
+ rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1126
netdata_systemd_journal_transform_gid, NULL);
1127
1128
facets_register_key_name_transformation(
1144
facets, "_SOURCE_REALTIME_TIMESTAMP",
1145
FACET_KEY_OPTION_TRANSFORM_VIEW,
1146
netdata_systemd_journal_transform_timestamp_usec, NULL);
1147
+}
1148
2095
- // ------------------------------------------------------------------------
2096
- // put this request into the progress db
2097
-
2098
- fqs = &tmp_fqs;
2099
-
2100
- // ------------------------------------------------------------------------
2101
- // validate parameters
2102
-
2103
- time_t now_s = now_realtime_sec();
2104
- time_t expires = now_s + 1;
2105
-
2106
- if(!q.after_s && !q.before_s) {
2107
- q.before_s = now_s;
2108
- q.after_s = q.before_s - SYSTEMD_JOURNAL_DEFAULT_QUERY_DURATION;
2109
- }
2110
- else
2111
- rrdr_relative_window_to_absolute(&q.after_s, &q.before_s, now_s);
2112
-
2113
- if(q.after_s > q.before_s) {
2114
- time_t tmp = q.after_s;
2115
- q.after_s = q.before_s;
2116
- q.before_s = tmp;
2117
- }
2118
-
2119
- if(q.after_s == q.before_s)
2120
- q.after_s = q.before_s - SYSTEMD_JOURNAL_DEFAULT_QUERY_DURATION;
2121
-
2122
- if(!q.last)
2123
- q.last = SYSTEMD_JOURNAL_DEFAULT_ITEMS_PER_QUERY;
2124
-
2125
- // ------------------------------------------------------------------------
2126
- // set query time-frame, anchors and direction
2127
-
2128
- fqs->transaction = transaction;
2129
- fqs->after_ut = q.after_s * USEC_PER_SEC;
2130
- fqs->before_ut = (q.before_s * USEC_PER_SEC) + USEC_PER_SEC - 1;
2131
- fqs->if_modified_since = q.if_modified_since;
2132
- fqs->data_only = q.data_only;
2133
- fqs->delta = (fqs->data_only) ? q.delta : false;
2134
- fqs->tail = (fqs->data_only && fqs->if_modified_since) ? q.tail : false;
2135
- fqs->sources = q.sources;
2136
- fqs->source_type = q.source_type;
2137
- fqs->entries = q.last;
2138
- fqs->last_modified = 0;
2139
- fqs->filters = q.filters;
2140
- fqs->query = (q.query && *q.query) ? q.query : NULL;
2141
- fqs->histogram = (q.chart && *q.chart) ? q.chart : NULL;
2142
- fqs->direction = q.direction;
2143
- fqs->anchor.start_ut = q.anchor;
2144
- fqs->anchor.stop_ut = 0;
2145
- fqs->sampling = q.sampling;
2146
-
2147
- if(fqs->anchor.start_ut && fqs->tail) {
2148
- // a tail request
2149
- // we need the top X entries from BEFORE
2150
- // but, we need to calculate the facets and the
2151
- // histogram up to the anchor
2152
- fqs->direction = q.direction = FACETS_ANCHOR_DIRECTION_BACKWARD;
2153
- fqs->anchor.start_ut = 0;
2154
- fqs->anchor.stop_ut = q.anchor;
2155
- }
2156
-
2157
- if(q.anchor && q.anchor < fqs->after_ut) {
2158
- log_fqs(fqs, "received anchor is too small for query timeframe, ignoring anchor");
2159
- q.anchor = 0;
2160
- fqs->anchor.start_ut = 0;
2161
- fqs->anchor.stop_ut = 0;
2162
- fqs->direction = q.direction = FACETS_ANCHOR_DIRECTION_BACKWARD;
2163
- }
2164
- else if(q.anchor > fqs->before_ut) {
2165
- log_fqs(fqs, "received anchor is too big for query timeframe, ignoring anchor");
2166
- q.anchor = 0;
2167
- fqs->anchor.start_ut = 0;
2168
- fqs->anchor.stop_ut = 0;
2169
- fqs->direction = q.direction = FACETS_ANCHOR_DIRECTION_BACKWARD;
2170
- }
2171
-
2172
- facets_set_anchor(facets, fqs->anchor.start_ut, fqs->anchor.stop_ut, fqs->direction);
2173
-
2174
- facets_set_additional_options(facets,
2175
- ((fqs->data_only) ? FACETS_OPTION_DATA_ONLY : 0) |
2176
- ((fqs->delta) ? FACETS_OPTION_SHOW_DELTAS : 0));
2177
-
2178
- // ------------------------------------------------------------------------
2179
- // set the rest of the query parameters
2180
-
2181
-
2182
- facets_set_items(facets, fqs->entries);
2183
- facets_set_query(facets, fqs->query);
1149
+void function_systemd_journal(const char *transaction, char *function, usec_t *stop_monotonic_ut, bool *cancelled,
1150
+ BUFFER *payload, HTTP_ACCESS access __maybe_unused,
1151
+ const char *source __maybe_unused, void *data __maybe_unused) {
1152
+ fstat_thread_calls = 0;
1153
+ fstat_thread_cached_responses = 0;
1154
1155
#ifdef HAVE_SD_JOURNAL_RESTART_FIELDS
2186
- fqs->slice = q.slice;
2187
- if(q.slice)
2188
- facets_enable_slice_mode(facets);
1156
+ bool have_slice = true;
1157
#else
2190
- fqs->slice = false;
2191
-#endif
2192
-
2193
- if(fqs->histogram)
2194
- facets_set_timeframe_and_histogram_by_id(facets, fqs->histogram, fqs->after_ut, fqs->before_ut);
2195
- else
2196
- facets_set_timeframe_and_histogram_by_name(facets, "PRIORITY", fqs->after_ut, fqs->before_ut);
1158
+ bool have_slice = false;
1159
+#endif // HAVE_SD_JOURNAL_RESTART_FIELDS
1160
1161
+ LOGS_QUERY_STATUS tmp_fqs = {
1162
+ .facets = lqs_facets_create(
1163
+ LQS_DEFAULT_ITEMS_PER_QUERY,
1164
+ FACETS_OPTION_ALL_KEYS_FTS,
1165
+ SYSTEMD_ALWAYS_VISIBLE_KEYS,
1166
+ SYSTEMD_KEYS_INCLUDED_IN_FACETS,
1167
+ SYSTEMD_KEYS_EXCLUDED_FROM_FACETS,
1168
+ have_slice),
1169
2199
- // ------------------------------------------------------------------------
2200
- // complete the request object
2201
-
2202
- buffer_json_member_add_boolean(wb, JOURNAL_PARAMETER_INFO, false);
2203
- buffer_json_member_add_boolean(wb, JOURNAL_PARAMETER_SLICE, fqs->slice);
2204
- buffer_json_member_add_boolean(wb, JOURNAL_PARAMETER_DATA_ONLY, fqs->data_only);
2205
- buffer_json_member_add_boolean(wb, JOURNAL_PARAMETER_DELTA, fqs->delta);
2206
- buffer_json_member_add_boolean(wb, JOURNAL_PARAMETER_TAIL, fqs->tail);
2207
- buffer_json_member_add_uint64(wb, JOURNAL_PARAMETER_SAMPLING, fqs->sampling);
2208
- buffer_json_member_add_uint64(wb, "source_type", fqs->source_type);
2209
- buffer_json_member_add_uint64(wb, JOURNAL_PARAMETER_AFTER, fqs->after_ut / USEC_PER_SEC);
2210
- buffer_json_member_add_uint64(wb, JOURNAL_PARAMETER_BEFORE, fqs->before_ut / USEC_PER_SEC);
2211
- buffer_json_member_add_uint64(wb, "if_modified_since", fqs->if_modified_since);
2212
- buffer_json_member_add_uint64(wb, JOURNAL_PARAMETER_ANCHOR, q.anchor);
2213
- buffer_json_member_add_string(wb, JOURNAL_PARAMETER_DIRECTION, fqs->direction == FACETS_ANCHOR_DIRECTION_FORWARD ? "forward" : "backward");
2214
- buffer_json_member_add_uint64(wb, JOURNAL_PARAMETER_LAST, fqs->entries);
2215
- buffer_json_member_add_string(wb, JOURNAL_PARAMETER_QUERY, fqs->query);
2216
- buffer_json_member_add_string(wb, JOURNAL_PARAMETER_HISTOGRAM, fqs->histogram);
2217
- buffer_json_object_close(wb); // request
2218
-
2219
- buffer_json_journal_versions(wb);
1170
+ .rq = LOGS_QUERY_REQUEST_DEFAULTS(transaction, LQS_DEFAULT_SLICE_MODE, JOURNAL_DEFAULT_DIRECTION),
1171
2221
- // ------------------------------------------------------------------------
2222
- // run the request
1172
+ .cancelled = cancelled,
1173
+ .stop_monotonic_ut = stop_monotonic_ut,
1174
+ };
1175
+ LOGS_QUERY_STATUS *lqs = &tmp_fqs;
1176
2224
- int response;
1177
+ CLEAN_BUFFER *wb = lqs_create_output_buffer();
1178
2226
- if(q.info) {
2227
- buffer_json_member_add_uint64(wb, "v", 3);
2228
- facets_accepted_parameters_to_json_array(facets, wb, false);
2229
- buffer_json_member_add_array(wb, "required_params");
2230
- {
2231
- buffer_json_add_array_item_object(wb);
2232
- {
2233
- buffer_json_member_add_string(wb, "id", "source");
2234
- buffer_json_member_add_string(wb, "name", "source");
2235
- buffer_json_member_add_string(wb, "help", "Select the SystemD Journal source to query");
2236
- buffer_json_member_add_string(wb, "type", "multiselect");
2237
- buffer_json_member_add_array(wb, "options");
2238
- {
2239
- available_journal_file_sources_to_json_array(wb);
2240
- }
2241
- buffer_json_array_close(wb); // options array
2242
- }
2243
- buffer_json_object_close(wb); // required params object
2244
- }
2245
- buffer_json_array_close(wb); // required_params array
1179
+ // ------------------------------------------------------------------------
1180
+ // parse the parameters
1181
2247
- facets_table_config(wb);
1182
+ if(lqs_request_parse_and_validate(lqs, wb, function, payload, have_slice, "PRIORITY")) {
1183
+ systemd_journal_register_transformations(lqs);
1184
2249
- buffer_json_member_add_uint64(wb, "status", HTTP_RESP_OK);
2250
- buffer_json_member_add_string(wb, "type", "table");
2251
- buffer_json_member_add_string(wb, "help", SYSTEMD_JOURNAL_FUNCTION_DESCRIPTION);
2252
- buffer_json_finalize(wb);
2253
- response = HTTP_RESP_OK;
2254
- goto output;
2255
- }
1185
+ // ------------------------------------------------------------------------
1186
+ // add versions to the response
1187
2257
- response = netdata_systemd_journal_query(wb, facets, fqs);
1188
+ buffer_json_journal_versions(wb);
1189
2259
- // ------------------------------------------------------------------------
2260
- // handle error response
1190
+ // ------------------------------------------------------------------------
1191
+ // run the request
1192
2262
- if(response != HTTP_RESP_OK) {
2263
- netdata_mutex_lock(&stdout_mutex);
2264
- pluginsd_function_json_error_to_stdout(transaction, response, "failed");
2265
- netdata_mutex_unlock(&stdout_mutex);
2266
- goto cleanup;
1193
+ if (lqs->rq.info)
1194
+ lqs_info_response(wb, lqs->facets);
1195
+ else {
1196
+ netdata_systemd_journal_query(wb, lqs);
1197
+ if (wb->response_code == HTTP_RESP_OK)
1198
+ buffer_json_finalize(wb);
1199
+ }
1200
}
1201
2269
-output:
1202
netdata_mutex_lock(&stdout_mutex);
2271
- pluginsd_function_result_to_stdout(transaction, response, "application/json", expires, wb);
1203
+ pluginsd_function_result_to_stdout(transaction, wb);
1204
netdata_mutex_unlock(&stdout_mutex);
1205
2274
-cleanup:
2275
- freez((void *)q.query);
2276
- freez((void *)q.chart);
2277
- simple_pattern_free(q.sources);
2278
- facets_destroy(facets);
1206
+ lqs_cleanup(lqs);
1207
}