1
+// SPDX-License-Identifier: GPL-3.0-or-later
2
+
3
+/*
4
+ * TODO
5
+ * _UDEV_DEVLINK is frequently set more than once per field - support multi-value faces
6
+ *
7
+ */
8
+
9
+#include "collectors/systemd-journal.plugin/provider/netdata_provider.h"
10
+#include "systemd-internals.h"
11
+
12
+#define ND_SD_JOURNAL_FUNCTION_DESCRIPTION "View, search and analyze systemd journal entries."
13
+#define ND_SD_JOURNAL_FUNCTION_NAME "systemd-journal"
14
+#define ND_SD_JOURNAL_SAMPLING_SLOTS 1000
15
+#define ND_SD_JOURNAL_SAMPLING_RECALIBRATE 10000
16
+
17
+#ifdef HAVE_SD_JOURNAL_RESTART_FIELDS
18
+#define LQS_DEFAULT_SLICE_MODE 1
19
+#else
20
+#define LQS_DEFAULT_SLICE_MODE 0
21
+#endif
22
+
23
+// functions needed by LQS
24
+static SD_JOURNAL_FILE_SOURCE_TYPE get_internal_source_type(const char *value);
25
+
26
+// structures needed by LQS
27
+struct lqs_extension {
28
+ struct {
29
+ usec_t start_ut;
30
+ usec_t stop_ut;
31
+ usec_t first_msg_ut;
32
+
33
+ NsdId128 first_msg_writer;
34
+ uint64_t first_msg_seqnum;
35
+ } query_file;
36
+
37
+ struct {
38
+ uint32_t enable_after_samples;
39
+ uint32_t slots;
40
+ uint32_t sampled;
41
+ uint32_t unsampled;
42
+ uint32_t estimated;
43
+ } samples;
44
+
45
+ struct {
46
+ uint32_t enable_after_samples;
47
+ uint32_t every;
48
+ uint32_t skipped;
49
+ uint32_t recalibrate;
50
+ uint32_t sampled;
51
+ uint32_t unsampled;
52
+ uint32_t estimated;
53
+ } samples_per_file;
54
+
55
+ struct {
56
+ usec_t start_ut;
57
+ usec_t end_ut;
58
+ usec_t step_ut;
59
+ uint32_t enable_after_samples;
60
+ uint32_t sampled[ND_SD_JOURNAL_SAMPLING_SLOTS];
61
+ uint32_t unsampled[ND_SD_JOURNAL_SAMPLING_SLOTS];
62
+ } samples_per_time_slot;
63
+
64
+ // per file progress info
65
+ // size_t cached_count;
66
+
67
+ // progress statistics
68
+ usec_t matches_setup_ut;
69
+ size_t rows_useful;
70
+ size_t rows_read;
71
+ size_t bytes_read;
72
+ size_t files_matched;
73
+ size_t file_working;
74
+};
75
+
76
+// prepare LQS
77
+#define LQS_FUNCTION_NAME ND_SD_JOURNAL_FUNCTION_NAME
78
+#define LQS_FUNCTION_DESCRIPTION ND_SD_JOURNAL_FUNCTION_DESCRIPTION
79
+#define LQS_DEFAULT_ITEMS_PER_QUERY 200
80
+#define LQS_DEFAULT_ITEMS_SAMPLING 1000000
81
+#define LQS_SOURCE_TYPE SD_JOURNAL_FILE_SOURCE_TYPE
82
+#define LQS_SOURCE_TYPE_ALL ND_SD_JF_ALL
83
+#define LQS_SOURCE_TYPE_NONE ND_SD_JF_NONE
84
+#define LQS_PARAMETER_SOURCE_NAME "Journal Sources" // this is how it is shown to users
85
+#define LQS_FUNCTION_GET_INTERNAL_SOURCE_TYPE(value) get_internal_source_type(value)
86
+#define LQS_FUNCTION_SOURCE_TO_JSON_ARRAY(wb) available_journal_file_sources_to_json_array(wb)
87
+#include "libnetdata/facets/logs_query_status.h"
88
+
89
+#include "systemd-journal-sampling.h"
90
+
91
+#define FACET_MAX_VALUE_LENGTH 8192
92
+#define ND_SD_JOURNAL_DEFAULT_TIMEOUT 60
93
+#define ND_SD_JOURNAL_PROGRESS_EVERY_UT (250 * USEC_PER_MS)
94
+#define JOURNAL_KEY_ND_JOURNAL_FILE "ND_JOURNAL_FILE"
95
+#define JOURNAL_KEY_ND_JOURNAL_PROCESS "ND_JOURNAL_PROCESS"
96
+#define JOURNAL_DEFAULT_DIRECTION FACETS_ANCHOR_DIRECTION_BACKWARD
97
+#define SYSTEMD_ALWAYS_VISIBLE_KEYS NULL
98
+
99
+#define SYSTEMD_KEYS_EXCLUDED_FROM_FACETS \
100
+ "!MESSAGE_ID" \
101
+ "|*MESSAGE*" \
102
+ "|*TIMESTAMP*" \
103
+ "|__*" \
104
+ ""
105
+
106
+#define SYSTEMD_KEYS_INCLUDED_IN_FACETS \
107
+ \
108
+ /* --- USER JOURNAL FIELDS --- */ \
109
+ \
110
+ /* "|MESSAGE" */ \
111
+ "|MESSAGE_ID" \
112
+ "|PRIORITY" \
113
+ "|CODE_FILE" /* "|CODE_LINE" */ \
114
+ "|CODE_FUNC" \
115
+ "|ERRNO" /* "|INVOCATION_ID" */ /* "|USER_INVOCATION_ID" */ \
116
+ "|SYSLOG_FACILITY" \
117
+ "|SYSLOG_IDENTIFIER" /* "|SYSLOG_PID" */ /* "|SYSLOG_TIMESTAMP" */ /* "|SYSLOG_RAW" */ /* "!DOCUMENTATION" */ /* "|TID" */ \
118
+ "|UNIT" \
119
+ "|USER_UNIT" \
120
+ "|UNIT_RESULT" /* undocumented */ \
121
+ \
122
+ /* --- TRUSTED JOURNAL FIELDS --- */ \
123
+ \
124
+ /* "|_PID" */ \
125
+ "|_UID" \
126
+ "|_GID" \
127
+ "|_COMM" \
128
+ "|_EXE" /* "|_CMDLINE" */ \
129
+ "|_CAP_EFFECTIVE" /* "|_AUDIT_SESSION" */ \
130
+ "|_AUDIT_LOGINUID" \
131
+ "|_SYSTEMD_CGROUP" \
132
+ "|_SYSTEMD_SLICE" \
133
+ "|_SYSTEMD_UNIT" \
134
+ "|_SYSTEMD_USER_UNIT" \
135
+ "|_SYSTEMD_USER_SLICE" \
136
+ "|_SYSTEMD_SESSION" \
137
+ "|_SYSTEMD_OWNER_UID" \
138
+ "|_SELINUX_CONTEXT" /* "|_SOURCE_REALTIME_TIMESTAMP" */ \
139
+ "|_BOOT_ID" \
140
+ "|_MACHINE_ID" /* "|_SYSTEMD_INVOCATION_ID" */ \
141
+ "|_HOSTNAME" \
142
+ "|_TRANSPORT" \
143
+ "|_STREAM_ID" /* "|LINE_BREAK" */ \
144
+ "|_NAMESPACE" \
145
+ "|_RUNTIME_SCOPE" \
146
+ \
147
+ /* --- KERNEL JOURNAL FIELDS --- */ \
148
+ \
149
+ /* "|_KERNEL_DEVICE" */ \
150
+ "|_KERNEL_SUBSYSTEM" /* "|_UDEV_SYSNAME" */ \
151
+ "|_UDEV_DEVNODE" /* "|_UDEV_DEVLINK" */ \
152
+ \
153
+ /* --- LOGGING ON BEHALF --- */ \
154
+ \
155
+ "|OBJECT_UID" \
156
+ "|OBJECT_GID" \
157
+ "|OBJECT_COMM" \
158
+ "|OBJECT_EXE" /* "|OBJECT_CMDLINE" */ /* "|OBJECT_AUDIT_SESSION" */ \
159
+ "|OBJECT_AUDIT_LOGINUID" \
160
+ "|OBJECT_SYSTEMD_CGROUP" \
161
+ "|OBJECT_SYSTEMD_SESSION" \
162
+ "|OBJECT_SYSTEMD_OWNER_UID" \
163
+ "|OBJECT_SYSTEMD_UNIT" \
164
+ "|OBJECT_SYSTEMD_USER_UNIT" \
165
+ \
166
+ /* --- CORE DUMPS --- */ \
167
+ \
168
+ "|COREDUMP_COMM" \
169
+ "|COREDUMP_UNIT" \
170
+ "|COREDUMP_USER_UNIT" \
171
+ "|COREDUMP_SIGNAL_NAME" \
172
+ "|COREDUMP_CGROUP" \
173
+ \
174
+ /* --- DOCKER --- */ \
175
+ \
176
+ "|CONTAINER_ID" /* "|CONTAINER_ID_FULL" */ \
177
+ "|CONTAINER_NAME" \
178
+ "|CONTAINER_TAG" \
179
+ "|IMAGE_NAME" /* undocumented */ /* "|CONTAINER_PARTIAL_MESSAGE" */ \
180
+ \
181
+ /* --- NETDATA --- */ \
182
+ \
183
+ "|ND_NIDL_NODE" \
184
+ "|ND_NIDL_CONTEXT" \
185
+ "|ND_LOG_SOURCE" /*"|ND_MODULE" */ \
186
+ "|ND_ALERT_NAME" \
187
+ "|ND_ALERT_CLASS" \
188
+ "|ND_ALERT_COMPONENT" \
189
+ "|ND_ALERT_TYPE" \
190
+ "|ND_ALERT_STATUS" \
191
+ \
192
+ ""
193
+
194
+static SD_JOURNAL_FILE_SOURCE_TYPE get_internal_source_type(const char *value)
195
+{
196
+ if (strcmp(value, ND_SD_JF_SOURCE_ALL_NAME) == 0)
197
+ return ND_SD_JF_ALL;
198
+ else if (strcmp(value, ND_SD_JF_SOURCE_LOCAL_NAME) == 0)
199
+ return ND_SD_JF_LOCAL_ALL;
200
+ else if (strcmp(value, ND_SD_JF_SOURCE_REMOTES_NAME) == 0)
201
+ return ND_SD_JF_REMOTE_ALL;
202
+ else if (strcmp(value, ND_SD_JF_SOURCE_NAMESPACES_NAME) == 0)
203
+ return ND_SD_JF_LOCAL_NAMESPACE;
204
+ else if (strcmp(value, ND_SD_JF_SOURCE_LOCAL_SYSTEM_NAME) == 0)
205
+ return ND_SD_JF_LOCAL_SYSTEM;
206
+ else if (strcmp(value, ND_SD_JF_SOURCE_LOCAL_USERS_NAME) == 0)
207
+ return ND_SD_JF_LOCAL_USER;
208
+ else if (strcmp(value, ND_SD_JF_SOURCE_LOCAL_OTHER_NAME) == 0)
209
+ return ND_SD_JF_LOCAL_OTHER;
210
+
211
+ return ND_SD_JF_NONE;
212
+}
213
+
214
+static inline bool nd_sd_journal_seek_to(NsdJournal *j, usec_t timestamp)
215
+{
216
+ if (nsd_journal_seek_realtime_usec(j, timestamp) < 0) {
217
+ netdata_log_error("SYSTEMD-JOURNAL: Failed to seek to %" PRIu64, timestamp);
218
+ if (nsd_journal_seek_tail(j) < 0) {
219
+ netdata_log_error("SYSTEMD-JOURNAL: Failed to seek to journal's tail");
220
+ return false;
221
+ }
222
+ }
223
+
224
+ return true;
225
+}
226
+
227
+#define JD_SOURCE_REALTIME_TIMESTAMP "_SOURCE_REALTIME_TIMESTAMP"
228
+
229
+static inline size_t
230
+nd_sd_journal_process_row(NsdJournal *j, FACETS *facets, struct nd_journal_file *njf, usec_t *msg_ut)
231
+{
232
+ const void *data;
233
+ size_t length, bytes = 0;
234
+
235
+ facets_add_key_value_length(
236
+ facets, JOURNAL_KEY_ND_JOURNAL_FILE, sizeof(JOURNAL_KEY_ND_JOURNAL_FILE) - 1, njf->filename, njf->filename_len);
237
+
238
+ NSD_JOURNAL_FOREACH_DATA(j, data, length)
239
+ {
240
+ const char *key, *value;
241
+ size_t key_length, value_length;
242
+
243
+ if (!parse_journal_field(data, length, &key, &key_length, &value, &value_length))
244
+ continue;
245
+
246
+#ifdef NETDATA_INTERNAL_CHECKS
247
+ usec_t origin_journal_ut = *msg_ut;
248
+#endif
249
+ if (unlikely(
250
+ key_length == sizeof(JD_SOURCE_REALTIME_TIMESTAMP) - 1 &&
251
+ memcmp(key, JD_SOURCE_REALTIME_TIMESTAMP, sizeof(JD_SOURCE_REALTIME_TIMESTAMP) - 1) == 0)) {
252
+ usec_t ut = str2ull(value, NULL);
253
+ if (ut && ut < *msg_ut) {
254
+ usec_t delta = *msg_ut - ut;
255
+ *msg_ut = ut;
256
+
257
+ if (delta > JOURNAL_VS_REALTIME_DELTA_MAX_UT)
258
+ delta = JOURNAL_VS_REALTIME_DELTA_MAX_UT;
259
+
260
+ // update max_journal_vs_realtime_delta_ut if the delta increased
261
+ usec_t expected = njf->max_journal_vs_realtime_delta_ut;
262
+ do {
263
+ if (delta <= expected)
264
+ break;
265
+ } while (!__atomic_compare_exchange_n(
266
+ &njf->max_journal_vs_realtime_delta_ut, &expected, delta, false, __ATOMIC_RELAXED, __ATOMIC_RELAXED));
267
+
268
+ internal_error(
269
+ delta > expected,
270
+ "increased max_journal_vs_realtime_delta_ut from %" PRIu64 " to %" PRIu64 ", "
271
+ "journal %" PRIu64 ", actual %" PRIu64 " (delta %" PRIu64 ")",
272
+ expected,
273
+ delta,
274
+ origin_journal_ut,
275
+ *msg_ut,
276
+ origin_journal_ut - (*msg_ut));
277
+ }
278
+ }
279
+
280
+ bytes += length;
281
+ facets_add_key_value_length(
282
+ facets,
283
+ key,
284
+ key_length,
285
+ value,
286
+ value_length <= FACET_MAX_VALUE_LENGTH ? value_length : FACET_MAX_VALUE_LENGTH);
287
+ }
288
+
289
+ return bytes;
290
+}
291
+
292
+#define FUNCTION_PROGRESS_UPDATE_ROWS(rows_read, rows) __atomic_fetch_add(&(rows_read), rows, __ATOMIC_RELAXED)
293
+#define FUNCTION_PROGRESS_UPDATE_BYTES(bytes_read, bytes) __atomic_fetch_add(&(bytes_read), bytes, __ATOMIC_RELAXED)
294
+#define FUNCTION_PROGRESS_EVERY_ROWS (1ULL << 13)
295
+#define FUNCTION_DATA_ONLY_CHECK_EVERY_ROWS (1ULL << 7)
296
+
297
+static inline ND_SD_JOURNAL_STATUS check_stop(const bool *cancelled, const usec_t *stop_monotonic_ut)
298
+{
299
+ if (cancelled && __atomic_load_n(cancelled, __ATOMIC_RELAXED)) {
300
+ internal_error(true, "Function has been cancelled");
301
+ return ND_SD_JOURNAL_CANCELLED;
302
+ }
303
+
304
+ if (now_monotonic_usec() > __atomic_load_n(stop_monotonic_ut, __ATOMIC_RELAXED)) {
305
+ internal_error(true, "Function timed out");
306
+ return ND_SD_JOURNAL_TIMED_OUT;
307
+ }
308
+
309
+ return ND_SD_JOURNAL_OK;
310
+}
311
+
312
+ND_SD_JOURNAL_STATUS nd_sd_journal_query_backward(
313
+ NsdJournal *j,
314
+ BUFFER *wb __maybe_unused,
315
+ FACETS *facets,
316
+ struct nd_journal_file *njf,
317
+ LOGS_QUERY_STATUS *fqs)
318
+{
319
+ usec_t anchor_delta = __atomic_load_n(&njf->max_journal_vs_realtime_delta_ut, __ATOMIC_RELAXED);
320
+ lqs_query_timeframe(fqs, anchor_delta);
321
+ usec_t start_ut = fqs->query.start_ut;
322
+ usec_t stop_ut = fqs->query.stop_ut;
323
+ bool stop_when_full = fqs->query.stop_when_full;
324
+
325
+ fqs->c.query_file.start_ut = start_ut;
326
+ fqs->c.query_file.stop_ut = stop_ut;
327
+
328
+ if (!nd_sd_journal_seek_to(j, start_ut))
329
+ return ND_SD_JOURNAL_FAILED_TO_SEEK;
330
+
331
+ size_t errors_no_timestamp = 0;
332
+ usec_t latest_msg_ut = 0; // the biggest timestamp we have seen so far
333
+ usec_t first_msg_ut = 0; // the first message we got from the db
334
+ size_t row_counter = 0, last_row_counter = 0, rows_useful = 0;
335
+ size_t bytes = 0, last_bytes = 0;
336
+
337
+ usec_t last_usec_from = 0;
338
+ usec_t last_usec_to = 0;
339
+
340
+ ND_SD_JOURNAL_STATUS status = ND_SD_JOURNAL_OK;
341
+
342
+ facets_rows_begin(facets);
343
+ while (status == ND_SD_JOURNAL_OK && nsd_journal_previous(j) > 0) {
344
+ usec_t msg_ut = 0;
345
+ if (nsd_journal_get_realtime_usec(j, &msg_ut) < 0 || !msg_ut) {
346
+ errors_no_timestamp++;
347
+ continue;
348
+ }
349
+
350
+ if (unlikely(msg_ut > start_ut))
351
+ continue;
352
+
353
+ if (unlikely(msg_ut < stop_ut))
354
+ break;
355
+
356
+ if (unlikely(msg_ut > latest_msg_ut))
357
+ latest_msg_ut = msg_ut;
358
+
359
+ if (unlikely(!first_msg_ut)) {
360
+ first_msg_ut = msg_ut;
361
+ fqs->c.query_file.first_msg_ut = msg_ut;
362
+
363
+#ifdef HAVE_SD_JOURNAL_GET_SEQNUM
364
+ if (nsd_journal_get_seqnum(j, &fqs->c.query_file.first_msg_seqnum, &fqs->c.query_file.first_msg_writer) <
365
+ 0) {
366
+ fqs->c.query_file.first_msg_seqnum = 0;
367
+ fqs->c.query_file.first_msg_writer = NSD_ID128_NULL;
368
+ }
369
+#endif
370
+ }
371
+
372
+ sampling_t sample = is_row_in_sample(
373
+ j, fqs, njf, msg_ut, FACETS_ANCHOR_DIRECTION_BACKWARD, facets_row_candidate_to_keep(facets, msg_ut));
374
+
375
+ if (sample == SAMPLING_FULL) {
376
+ bytes += nd_sd_journal_process_row(j, facets, njf, &msg_ut);
377
+
378
+ // make sure each line gets a unique timestamp
379
+ if (unlikely(msg_ut >= last_usec_from && msg_ut <= last_usec_to))
380
+ msg_ut = --last_usec_from;
381
+ else
382
+ last_usec_from = last_usec_to = msg_ut;
383
+
384
+ if (facets_row_finished(facets, msg_ut))
385
+ rows_useful++;
386
+
387
+ row_counter++;
388
+ if (unlikely(
389
+ (row_counter % FUNCTION_DATA_ONLY_CHECK_EVERY_ROWS) == 0 && stop_when_full &&
390
+ facets_rows(facets) >= fqs->rq.entries)) {
391
+ // stop the data only query
392
+ usec_t oldest = facets_row_oldest_ut(facets);
393
+ if (oldest && msg_ut < (oldest - anchor_delta))
394
+ break;
395
+ }
396
+
397
+ if (unlikely(row_counter % FUNCTION_PROGRESS_EVERY_ROWS == 0)) {
398
+ FUNCTION_PROGRESS_UPDATE_ROWS(fqs->c.rows_read, row_counter - last_row_counter);
399
+ last_row_counter = row_counter;
400
+
401
+ FUNCTION_PROGRESS_UPDATE_BYTES(fqs->c.bytes_read, bytes - last_bytes);
402
+ last_bytes = bytes;
403
+
404
+ status = check_stop(fqs->cancelled, fqs->stop_monotonic_ut);
405
+ }
406
+ } else if (sample == SAMPLING_SKIP_FIELDS)
407
+ facets_row_finished_unsampled(facets, msg_ut);
408
+ else {
409
+ sampling_update_running_query_file_estimates(facets, j, fqs, njf, msg_ut, FACETS_ANCHOR_DIRECTION_BACKWARD);
410
+ break;
411
+ }
412
+ }
413
+
414
+ FUNCTION_PROGRESS_UPDATE_ROWS(fqs->c.rows_read, row_counter - last_row_counter);
415
+ FUNCTION_PROGRESS_UPDATE_BYTES(fqs->c.bytes_read, bytes - last_bytes);
416
+
417
+ fqs->c.rows_useful += rows_useful;
418
+
419
+ if (errors_no_timestamp)
420
+ netdata_log_error("SYSTEMD-JOURNAL: %zu lines did not have timestamps", errors_no_timestamp);
421
+
422
+ if (latest_msg_ut > fqs->last_modified)
423
+ fqs->last_modified = latest_msg_ut;
424
+
425
+ return status;
426
+}
427
+
428
+ND_SD_JOURNAL_STATUS nd_sd_journal_query_forward(
429
+ NsdJournal *j,
430
+ BUFFER *wb __maybe_unused,
431
+ FACETS *facets,
432
+ struct nd_journal_file *njf,
433
+ LOGS_QUERY_STATUS *fqs)
434
+{
435
+ usec_t anchor_delta = __atomic_load_n(&njf->max_journal_vs_realtime_delta_ut, __ATOMIC_RELAXED);
436
+ lqs_query_timeframe(fqs, anchor_delta);
437
+ usec_t start_ut = fqs->query.start_ut;
438
+ usec_t stop_ut = fqs->query.stop_ut;
439
+ bool stop_when_full = fqs->query.stop_when_full;
440
+
441
+ fqs->c.query_file.start_ut = start_ut;
442
+ fqs->c.query_file.stop_ut = stop_ut;
443
+
444
+ if (!nd_sd_journal_seek_to(j, start_ut))
445
+ return ND_SD_JOURNAL_FAILED_TO_SEEK;
446
+
447
+ size_t errors_no_timestamp = 0;
448
+ usec_t latest_msg_ut = 0; // the biggest timestamp we have seen so far
449
+ usec_t first_msg_ut = 0; // the first message we got from the db
450
+ size_t row_counter = 0, last_row_counter = 0, rows_useful = 0;
451
+ size_t bytes = 0, last_bytes = 0;
452
+
453
+ usec_t last_usec_from = 0;
454
+ usec_t last_usec_to = 0;
455
+
456
+ ND_SD_JOURNAL_STATUS status = ND_SD_JOURNAL_OK;
457
+
458
+ facets_rows_begin(facets);
459
+ while (status == ND_SD_JOURNAL_OK && nsd_journal_next(j) > 0) {
460
+ usec_t msg_ut = 0;
461
+ if (nsd_journal_get_realtime_usec(j, &msg_ut) < 0 || !msg_ut) {
462
+ errors_no_timestamp++;
463
+ continue;
464
+ }
465
+
466
+ if (unlikely(msg_ut < start_ut))
467
+ continue;
468
+
469
+ if (unlikely(msg_ut > stop_ut))
470
+ break;
471
+
472
+ if (likely(msg_ut > latest_msg_ut))
473
+ latest_msg_ut = msg_ut;
474
+
475
+ if (unlikely(!first_msg_ut)) {
476
+ first_msg_ut = msg_ut;
477
+ fqs->c.query_file.first_msg_ut = msg_ut;
478
+ }
479
+
480
+ sampling_t sample = is_row_in_sample(
481
+ j, fqs, njf, msg_ut, FACETS_ANCHOR_DIRECTION_FORWARD, facets_row_candidate_to_keep(facets, msg_ut));
482
+
483
+ if (sample == SAMPLING_FULL) {
484
+ bytes += nd_sd_journal_process_row(j, facets, njf, &msg_ut);
485
+
486
+ // make sure each line gets a unique timestamp
487
+ if (unlikely(msg_ut >= last_usec_from && msg_ut <= last_usec_to))
488
+ msg_ut = ++last_usec_to;
489
+ else
490
+ last_usec_from = last_usec_to = msg_ut;
491
+
492
+ if (facets_row_finished(facets, msg_ut))
493
+ rows_useful++;
494
+
495
+ row_counter++;
496
+ if (unlikely(
497
+ (row_counter % FUNCTION_DATA_ONLY_CHECK_EVERY_ROWS) == 0 && stop_when_full &&
498
+ facets_rows(facets) >= fqs->rq.entries)) {
499
+ // stop the data only query
500
+ usec_t newest = facets_row_newest_ut(facets);
501
+ if (newest && msg_ut > (newest + anchor_delta))
502
+ break;
503
+ }
504
+
505
+ if (unlikely(row_counter % FUNCTION_PROGRESS_EVERY_ROWS == 0)) {
506
+ FUNCTION_PROGRESS_UPDATE_ROWS(fqs->c.rows_read, row_counter - last_row_counter);
507
+ last_row_counter = row_counter;
508
+
509
+ FUNCTION_PROGRESS_UPDATE_BYTES(fqs->c.bytes_read, bytes - last_bytes);
510
+ last_bytes = bytes;
511
+
512
+ status = check_stop(fqs->cancelled, fqs->stop_monotonic_ut);
513
+ }
514
+ } else if (sample == SAMPLING_SKIP_FIELDS)
515
+ facets_row_finished_unsampled(facets, msg_ut);
516
+ else {
517
+ sampling_update_running_query_file_estimates(facets, j, fqs, njf, msg_ut, FACETS_ANCHOR_DIRECTION_FORWARD);
518
+ break;
519
+ }
520
+ }
521
+
522
+ FUNCTION_PROGRESS_UPDATE_ROWS(fqs->c.rows_read, row_counter - last_row_counter);
523
+ FUNCTION_PROGRESS_UPDATE_BYTES(fqs->c.bytes_read, bytes - last_bytes);
524
+
525
+ fqs->c.rows_useful += rows_useful;
526
+
527
+ if (errors_no_timestamp)
528
+ netdata_log_error("SYSTEMD-JOURNAL: %zu lines did not have timestamps", errors_no_timestamp);
529
+
530
+ if (latest_msg_ut > fqs->last_modified)
531
+ fqs->last_modified = latest_msg_ut;
532
+
533
+ return status;
534
+}
535
+
536
+bool nd_sd_journal_check_if_modified_since(NsdJournal *j, usec_t seek_to, usec_t last_modified)
537
+{
538
+ // return true, if data have been modified since the timestamp
539
+
540
+ if (!last_modified || !seek_to)
541
+ return false;
542
+
543
+ if (!nd_sd_journal_seek_to(j, seek_to))
544
+ return false;
545
+
546
+ usec_t first_msg_ut = 0;
547
+ while (nsd_journal_previous(j) > 0) {
548
+ usec_t msg_ut;
549
+ if (nsd_journal_get_realtime_usec(j, &msg_ut) < 0)
550
+ continue;
551
+
552
+ first_msg_ut = msg_ut;
553
+ break;
554
+ }
555
+
556
+ return first_msg_ut != last_modified;
557
+}
558
+
559
+#ifdef HAVE_SD_JOURNAL_RESTART_FIELDS
560
+static bool netdata_systemd_filtering_by_journal(NsdJournal *j, FACETS *facets, LOGS_QUERY_STATUS *lqs)
561
+{
562
+ const char *field = NULL;
563
+ const void *data = NULL;
564
+ size_t data_length;
565
+ size_t added_keys = 0;
566
+ size_t failures = 0;
567
+ size_t filters_added = 0;
568
+
569
+ NSD_JOURNAL_FOREACH_FIELD(j, field)
570
+ { // for each key
571
+ bool interesting;
572
+
573
+ if (lqs->rq.data_only)
574
+ interesting = facets_key_name_is_filter(facets, field);
575
+ else
576
+ interesting = facets_key_name_is_facet(facets, field);
577
+
578
+ if (interesting) {
579
+ if (nsd_journal_query_unique(j, field) >= 0) {
580
+ bool added_this_key = false;
581
+ size_t added_values = 0;
582
+
583
+ NSD_JOURNAL_FOREACH_UNIQUE(j, data, data_length)
584
+ { // for each value of the key
585
+ const char *key, *value;
586
+ size_t key_length, value_length;
587
+
588
+ if (!parse_journal_field(data, data_length, &key, &key_length, &value, &value_length))
589
+ continue;
590
+
591
+ facets_add_possible_value_name_to_key(facets, key, key_length, value, value_length);
592
+
593
+ if (!facets_key_name_value_length_is_selected(facets, key, key_length, value, value_length))
594
+ continue;
595
+
596
+ if (added_keys && !added_this_key) {
597
+ if (nsd_journal_add_conjunction(j) < 0) // key AND key AND key
598
+ failures++;
599
+
600
+ added_this_key = true;
601
+ added_keys++;
602
+ } else if (added_values)
603
+ if (nsd_journal_add_disjunction(j) < 0) // value OR value OR value
604
+ failures++;
605
+
606
+ if (nsd_journal_add_match(j, data, data_length) < 0)
607
+ failures++;
608
+
609
+ if (!added_keys) {
610
+ added_keys++;
611
+ added_this_key = true;
612
+ }
613
+
614
+ added_values++;
615
+ filters_added++;
616
+ }
617
+ }
618
+ }
619
+ }
620
+
621
+ if (failures) {
622
+ lqs_log_error(lqs, "failed to setup journal filter, will run the full query.");
623
+ nsd_journal_flush_matches(j);
624
+ return true;
625
+ }
626
+
627
+ return filters_added ? true : false;
628
+}
629
+#endif // HAVE_SD_JOURNAL_RESTART_FIELDS
630
+
631
+static ND_SD_JOURNAL_STATUS nd_sd_journal_query_one_file(
632
+ const char *filename,
633
+ BUFFER *wb,
634
+ FACETS *facets,
635
+ struct nd_journal_file *njf,
636
+ LOGS_QUERY_STATUS *fqs)
637
+{
638
+ NsdJournal *j = NULL;
639
+ errno_clear();
640
+
641
+ fstat_cache_enable_on_thread();
642
+
643
+ const char *paths[2] = {
644
+ [0] = filename,
645
+ [1] = NULL,
646
+ };
647
+
648
+ if (nsd_journal_open_files(&j, paths, ND_SD_JOURNAL_OPEN_FLAGS) < 0 || !j) {
649
+ netdata_log_error("JOURNAL: cannot open file '%s' for query", filename);
650
+ fstat_cache_disable_on_thread();
651
+ return ND_SD_JOURNAL_FAILED_TO_OPEN;
652
+ }
653
+
654
+ ND_SD_JOURNAL_STATUS status;
655
+ bool matches_filters = true;
656
+
657
+#ifdef HAVE_SD_JOURNAL_RESTART_FIELDS
658
+ if (fqs->rq.slice) {
659
+ usec_t started = now_monotonic_usec();
660
+
661
+ matches_filters = netdata_systemd_filtering_by_journal(j, facets, fqs) || !fqs->rq.filters;
662
+ usec_t ended = now_monotonic_usec();
663
+
664
+ fqs->c.matches_setup_ut += (ended - started);
665
+ }
666
+#endif // HAVE_SD_JOURNAL_RESTART_FIELDS
667
+
668
+ if (matches_filters) {
669
+ if (fqs->rq.direction == FACETS_ANCHOR_DIRECTION_FORWARD)
670
+ status = nd_sd_journal_query_forward(j, wb, facets, njf, fqs);
671
+ else
672
+ status = nd_sd_journal_query_backward(j, wb, facets, njf, fqs);
673
+ } else
674
+ status = ND_SD_JOURNAL_NO_FILE_MATCHED;
675
+
676
+ nsd_journal_close(j);
677
+ fstat_cache_disable_on_thread();
678
+
679
+ return status;
680
+}
681
+
682
+static bool jf_is_mine(struct nd_journal_file *njf, LOGS_QUERY_STATUS *fqs)
683
+{
684
+ if ((fqs->rq.source_type == ND_SD_JF_NONE && !fqs->rq.sources) || (njf->source_type & fqs->rq.source_type) ||
685
+ (fqs->rq.sources && simple_pattern_matches(fqs->rq.sources, string2str(njf->source)))) {
686
+ if (!njf->msg_last_ut)
687
+ // the file is not scanned yet, or the timestamps have not been updated,
688
+ // so we don't know if it can contribute or not - let's add it.
689
+ return true;
690
+
691
+ usec_t anchor_delta = JOURNAL_VS_REALTIME_DELTA_MAX_UT;
692
+ usec_t first_ut = njf->msg_first_ut - anchor_delta;
693
+ usec_t last_ut = njf->msg_last_ut + anchor_delta;
694
+
695
+ if (last_ut >= fqs->rq.after_ut && first_ut <= fqs->rq.before_ut)
696
+ return true;
697
+ }
698
+
699
+ return false;
700
+}
701
+
702
+static int nd_sd_journal_query(BUFFER *wb, LOGS_QUERY_STATUS *lqs)
703
+{
704
+ FACETS *facets = lqs->facets;
705
+
706
+ ND_SD_JOURNAL_STATUS status = ND_SD_JOURNAL_NO_FILE_MATCHED;
707
+ struct nd_journal_file *njf;
708
+
709
+ lqs->c.files_matched = 0;
710
+ lqs->c.file_working = 0;
711
+ lqs->c.rows_useful = 0;
712
+ lqs->c.rows_read = 0;
713
+ lqs->c.bytes_read = 0;
714
+
715
+ size_t files_used = 0;
716
+ size_t files_max = dictionary_entries(nd_journal_files_registry);
717
+ const DICTIONARY_ITEM *file_items[files_max];
718
+
719
+ // count the files
720
+ bool files_are_newer = false;
721
+ dfe_start_read(nd_journal_files_registry, njf)
722
+ {
723
+ if (!jf_is_mine(njf, lqs))
724
+ continue;
725
+
726
+ file_items[files_used++] = dictionary_acquired_item_dup(nd_journal_files_registry, njf_dfe.item);
727
+
728
+ if (njf->msg_last_ut > lqs->rq.if_modified_since)
729
+ files_are_newer = true;
730
+ }
731
+ dfe_done(jf);
732
+
733
+ lqs->c.files_matched = files_used;
734
+
735
+ if (lqs->rq.if_modified_since && !files_are_newer) {
736
+ // release the files
737
+ for (size_t f = 0; f < files_used; f++)
738
+ dictionary_acquired_item_release(nd_journal_files_registry, file_items[f]);
739
+
740
+ return rrd_call_function_error(wb, "No new data since the previous call.", HTTP_RESP_NOT_MODIFIED);
741
+ }
742
+
743
+ // sort the files, so that they are optimal for facets
744
+ if (files_used >= 2) {
745
+ if (lqs->rq.direction == FACETS_ANCHOR_DIRECTION_BACKWARD)
746
+ qsort(file_items, files_used, sizeof(const DICTIONARY_ITEM *), nd_journal_file_dict_items_backward_compar);
747
+ else
748
+ qsort(file_items, files_used, sizeof(const DICTIONARY_ITEM *), nd_journal_file_dict_items_forward_compar);
749
+ }
750
+
751
+ bool partial = false;
752
+ usec_t query_started_ut = now_monotonic_usec();
753
+ usec_t started_ut = query_started_ut;
754
+ usec_t ended_ut = started_ut;
755
+ usec_t duration_ut = 0, max_duration_ut = 0;
756
+ usec_t progress_duration_ut = 0;
757
+
758
+ sampling_query_init(lqs, facets);
759
+
760
+ buffer_json_member_add_array(wb, "_journal_files");
761
+ for (size_t f = 0; f < files_used; f++) {
762
+ const char *filename = dictionary_acquired_item_name(file_items[f]);
763
+ njf = dictionary_acquired_item_value(file_items[f]);
764
+
765
+ if (!jf_is_mine(njf, lqs))
766
+ continue;
767
+
768
+ started_ut = ended_ut;
769
+
770
+ // do not even try to do the query if we expect it to pass the timeout
771
+ if (ended_ut + max_duration_ut * 3 >= *lqs->stop_monotonic_ut) {
772
+ partial = true;
773
+ status = ND_SD_JOURNAL_TIMED_OUT;
774
+ break;
775
+ }
776
+
777
+ lqs->c.file_working++;
778
+
779
+ size_t fs_calls = fstat_thread_calls;
780
+ size_t fs_cached = fstat_thread_cached_responses;
781
+ size_t rows_useful = lqs->c.rows_useful;
782
+ size_t rows_read = lqs->c.rows_read;
783
+ size_t bytes_read = lqs->c.bytes_read;
784
+ size_t matches_setup_ut = lqs->c.matches_setup_ut;
785
+
786
+ sampling_file_init(lqs, njf);
787
+
788
+ ND_SD_JOURNAL_STATUS tmp_status = nd_sd_journal_query_one_file(filename, wb, facets, njf, lqs);
789
+
790
+ rows_useful = lqs->c.rows_useful - rows_useful;
791
+ rows_read = lqs->c.rows_read - rows_read;
792
+ bytes_read = lqs->c.bytes_read - bytes_read;
793
+ matches_setup_ut = lqs->c.matches_setup_ut - matches_setup_ut;
794
+ fs_calls = fstat_thread_calls - fs_calls;
795
+ fs_cached = fstat_thread_cached_responses - fs_cached;
796
+
797
+ ended_ut = now_monotonic_usec();
798
+ duration_ut = ended_ut - started_ut;
799
+
800
+ if (duration_ut > max_duration_ut)
801
+ max_duration_ut = duration_ut;
802
+
803
+ progress_duration_ut += duration_ut;
804
+ if (progress_duration_ut >= ND_SD_JOURNAL_PROGRESS_EVERY_UT) {
805
+ progress_duration_ut = 0;
806
+ netdata_mutex_lock(&stdout_mutex);
807
+ pluginsd_function_progress_to_stdout(lqs->rq.transaction, f + 1, files_used);
808
+ netdata_mutex_unlock(&stdout_mutex);
809
+ }
810
+
811
+ buffer_json_add_array_item_object(wb); // journal file
812
+ {
813
+ // information about the file
814
+ buffer_json_member_add_string(wb, "_filename", filename);
815
+ buffer_json_member_add_uint64(wb, "_source_type", njf->source_type);
816
+ buffer_json_member_add_string(wb, "_source", string2str(njf->source));
817
+ buffer_json_member_add_uint64(wb, "_last_modified_ut", njf->file_last_modified_ut);
818
+ buffer_json_member_add_uint64(wb, "_msg_first_ut", njf->msg_first_ut);
819
+ buffer_json_member_add_uint64(wb, "_msg_last_ut", njf->msg_last_ut);
820
+ buffer_json_member_add_uint64(wb, "_journal_vs_realtime_delta_ut", njf->max_journal_vs_realtime_delta_ut);
821
+
822
+ // information about the current use of the file
823
+ buffer_json_member_add_uint64(wb, "duration_ut", ended_ut - started_ut);
824
+ buffer_json_member_add_uint64(wb, "rows_read", rows_read);
825
+ buffer_json_member_add_uint64(wb, "rows_useful", rows_useful);
826
+ buffer_json_member_add_double(
827
+ wb, "rows_per_second", (double)rows_read / (double)duration_ut * (double)USEC_PER_SEC);
828
+ buffer_json_member_add_uint64(wb, "bytes_read", bytes_read);
829
+ buffer_json_member_add_double(
830
+ wb, "bytes_per_second", (double)bytes_read / (double)duration_ut * (double)USEC_PER_SEC);
831
+ buffer_json_member_add_uint64(wb, "duration_matches_ut", matches_setup_ut);
832
+ buffer_json_member_add_uint64(wb, "fstat_query_calls", fs_calls);
833
+ buffer_json_member_add_uint64(wb, "fstat_query_cached_responses", fs_cached);
834
+
835
+ if (lqs->rq.sampling) {
836
+ buffer_json_member_add_object(wb, "_sampling");
837
+ {
838
+ buffer_json_member_add_uint64(wb, "sampled", lqs->c.samples_per_file.sampled);
839
+ buffer_json_member_add_uint64(wb, "unsampled", lqs->c.samples_per_file.unsampled);
840
+ buffer_json_member_add_uint64(wb, "estimated", lqs->c.samples_per_file.estimated);
841
+ }
842
+ buffer_json_object_close(wb); // _sampling
843
+ }
844
+ }
845
+ buffer_json_object_close(wb); // journal file
846
+
847
+ bool stop = false;
848
+ switch (tmp_status) {
849
+ case ND_SD_JOURNAL_OK:
850
+ case ND_SD_JOURNAL_NO_FILE_MATCHED:
851
+ status = (status == ND_SD_JOURNAL_OK) ? ND_SD_JOURNAL_OK : tmp_status;
852
+ break;
853
+
854
+ case ND_SD_JOURNAL_FAILED_TO_OPEN:
855
+ case ND_SD_JOURNAL_FAILED_TO_SEEK:
856
+ partial = true;
857
+ if (status == ND_SD_JOURNAL_NO_FILE_MATCHED)
858
+ status = tmp_status;
859
+ break;
860
+
861
+ case ND_SD_JOURNAL_CANCELLED:
862
+ case ND_SD_JOURNAL_TIMED_OUT:
863
+ partial = true;
864
+ stop = true;
865
+ status = tmp_status;
866
+ break;
867
+
868
+ case ND_SD_JOURNAL_NOT_MODIFIED:
869
+ internal_fatal(true, "this should never be returned here");
870
+ break;
871
+ }
872
+
873
+ if (stop)
874
+ break;
875
+ }
876
+ buffer_json_array_close(wb); // _journal_files
877
+
878
+ // release the files
879
+ for (size_t f = 0; f < files_used; f++)
880
+ dictionary_acquired_item_release(nd_journal_files_registry, file_items[f]);
881
+
882
+ switch (status) {
883
+ case ND_SD_JOURNAL_OK:
884
+ if (lqs->rq.if_modified_since && !lqs->c.rows_useful)
885
+ return rrd_call_function_error(
886
+ wb, "No additional useful data since the previous call.", HTTP_RESP_NOT_MODIFIED);
887
+ break;
888
+
889
+ case ND_SD_JOURNAL_TIMED_OUT:
890
+ case ND_SD_JOURNAL_NO_FILE_MATCHED:
891
+ break;
892
+
893
+ case ND_SD_JOURNAL_CANCELLED:
894
+ return rrd_call_function_error(wb, "Request cancelled.", HTTP_RESP_CLIENT_CLOSED_REQUEST);
895
+
896
+ case ND_SD_JOURNAL_NOT_MODIFIED:
897
+ return rrd_call_function_error(wb, "No new data since the previous call.", HTTP_RESP_NOT_MODIFIED);
898
+
899
+ case ND_SD_JOURNAL_FAILED_TO_OPEN:
900
+ return rrd_call_function_error(wb, "Failed to open systemd journal file.", HTTP_RESP_INTERNAL_SERVER_ERROR);
901
+
902
+ case ND_SD_JOURNAL_FAILED_TO_SEEK:
903
+ return rrd_call_function_error(
904
+ wb, "Failed to seek in systemd journal file.", HTTP_RESP_INTERNAL_SERVER_ERROR);
905
+
906
+ default:
907
+ return rrd_call_function_error(wb, "Unknown status", HTTP_RESP_INTERNAL_SERVER_ERROR);
908
+ }
909
+
910
+ buffer_json_member_add_uint64(wb, "status", HTTP_RESP_OK);
911
+ buffer_json_member_add_boolean(wb, "partial", partial);
912
+ buffer_json_member_add_string(wb, "type", "table");
913
+
914
+ // build a message for the query
915
+ if (!lqs->rq.data_only) {
916
+ CLEAN_BUFFER *msg = buffer_create(0, NULL);
917
+ CLEAN_BUFFER *msg_description = buffer_create(0, NULL);
918
+ ND_LOG_FIELD_PRIORITY msg_priority = NDLP_INFO;
919
+
920
+ if (!nd_journal_files_completed_once()) {
921
+ buffer_strcat(msg, "Journals are still being scanned. ");
922
+ buffer_strcat(
923
+ msg_description,
924
+ "LIBRARY SCAN: The journal files are still being scanned, you are probably viewing incomplete data. ");
925
+ msg_priority = NDLP_WARNING;
926
+ }
927
+
928
+ if (partial) {
929
+ buffer_strcat(msg, "Query timed-out, incomplete data. ");
930
+ buffer_strcat(
931
+ msg_description,
932
+ "QUERY TIMEOUT: The query timed out and may not include all the data of the selected window. ");
933
+ msg_priority = NDLP_WARNING;
934
+ }
935
+
936
+ if (lqs->c.samples.estimated || lqs->c.samples.unsampled) {
937
+ double percent = (double)(lqs->c.samples.sampled * 100.0 /
938
+ (lqs->c.samples.estimated + lqs->c.samples.unsampled + lqs->c.samples.sampled));
939
+ buffer_sprintf(msg, "%.2f%% real data", percent);
940
+ buffer_sprintf(msg_description, "ACTUAL DATA: The filters counters reflect %0.2f%% of the data. ", percent);
941
+ msg_priority = MIN(msg_priority, NDLP_NOTICE);
942
+ }
943
+
944
+ if (lqs->c.samples.unsampled) {
945
+ double percent = (double)(lqs->c.samples.unsampled * 100.0 /
946
+ (lqs->c.samples.estimated + lqs->c.samples.unsampled + lqs->c.samples.sampled));
947
+ buffer_sprintf(msg, ", %.2f%% unsampled", percent);
948
+ buffer_sprintf(
949
+ msg_description,
950
+ "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. ",
951
+ percent);
952
+ msg_priority = MIN(msg_priority, NDLP_NOTICE);
953
+ }
954
+
955
+ if (lqs->c.samples.estimated) {
956
+ double percent = (double)(lqs->c.samples.estimated * 100.0 /
957
+ (lqs->c.samples.estimated + lqs->c.samples.unsampled + lqs->c.samples.sampled));
958
+ buffer_sprintf(msg, ", %.2f%% estimated", percent);
959
+ buffer_sprintf(
960
+ msg_description,
961
+ "ESTIMATED DATA: The query selected a large amount of data, so to avoid delaying too much, the presented data are estimated by %0.2f%%. ",
962
+ percent);
963
+ msg_priority = MIN(msg_priority, NDLP_NOTICE);
964
+ }
965
+
966
+ buffer_json_member_add_object(wb, "message");
967
+ if (buffer_tostring(msg)) {
968
+ buffer_json_member_add_string(wb, "title", buffer_tostring(msg));
969
+ buffer_json_member_add_string(wb, "description", buffer_tostring(msg_description));
970
+ buffer_json_member_add_string(wb, "status", nd_log_id2priority(msg_priority));
971
+ }
972
+ // else send an empty object if there is nothing to tell
973
+ buffer_json_object_close(wb); // message
974
+ }
975
+
976
+ if (!lqs->rq.data_only) {
977
+ buffer_json_member_add_time_t(wb, "update_every", 1);
978
+ buffer_json_member_add_string(wb, "help", ND_SD_JOURNAL_FUNCTION_DESCRIPTION);
979
+ }
980
+
981
+ if (!lqs->rq.data_only || lqs->rq.tail)
982
+ buffer_json_member_add_uint64(wb, "last_modified", lqs->last_modified);
983
+
984
+ facets_sort_and_reorder_keys(facets);
985
+ facets_report(facets, wb, used_hashes_registry);
986
+
987
+ wb->expires = now_realtime_sec() + (lqs->rq.data_only ? 3600 : 0);
988
+ buffer_json_member_add_time_t(wb, "expires", wb->expires);
989
+
990
+ buffer_json_member_add_object(wb, "_fstat_caching");
991
+ {
992
+ buffer_json_member_add_uint64(wb, "calls", fstat_thread_calls);
993
+ buffer_json_member_add_uint64(wb, "cached", fstat_thread_cached_responses);
994
+ }
995
+ buffer_json_object_close(wb); // _fstat_caching
996
+
997
+ if (lqs->rq.sampling) {
998
+ buffer_json_member_add_object(wb, "_sampling");
999
+ {
1000
+ buffer_json_member_add_uint64(wb, "sampled", lqs->c.samples.sampled);
1001
+ buffer_json_member_add_uint64(wb, "unsampled", lqs->c.samples.unsampled);
1002
+ buffer_json_member_add_uint64(wb, "estimated", lqs->c.samples.estimated);
1003
+ }
1004
+ buffer_json_object_close(wb); // _sampling
1005
+ }
1006
+
1007
+ wb->content_type = CT_APPLICATION_JSON;
1008
+ wb->response_code = HTTP_RESP_OK;
1009
+ return wb->response_code;
1010
+}
1011
+
1012
+static void systemd_journal_register_transformations(LOGS_QUERY_STATUS *lqs)
1013
+{
1014
+ FACETS *facets = lqs->facets;
1015
+ LOGS_QUERY_REQUEST *rq = &lqs->rq;
1016
+
1017
+ // ----------------------------------------------------------------------------------------------------------------
1018
+ // register the fields in the order you want them on the dashboard
1019
+
1020
+ facets_register_row_severity(facets, syslog_priority_to_facet_severity, NULL);
1021
+
1022
+ facets_register_key_name(facets, "_HOSTNAME", rq->default_facet | FACET_KEY_OPTION_VISIBLE);
1023
+
1024
+ facets_register_dynamic_key_name(
1025
+ facets,
1026
+ JOURNAL_KEY_ND_JOURNAL_PROCESS,
1027
+ FACET_KEY_OPTION_NEVER_FACET | FACET_KEY_OPTION_VISIBLE,
1028
+ nd_sd_journal_dynamic_row_id,
1029
+ NULL);
1030
+
1031
+ facets_register_key_name(
1032
+ facets,
1033
+ "MESSAGE",
1034
+ FACET_KEY_OPTION_NEVER_FACET | FACET_KEY_OPTION_MAIN_TEXT | FACET_KEY_OPTION_VISIBLE | FACET_KEY_OPTION_FTS);
1035
+
1036
+ facets_register_key_name_transformation(
1037
+ facets,
1038
+ "PRIORITY",
1039
+ rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW | FACET_KEY_OPTION_EXPANDED_FILTER,
1040
+ nd_sd_journal_transform_priority,
1041
+ NULL);
1042
+
1043
+ facets_register_key_name_transformation(
1044
+ facets,
1045
+ "SYSLOG_FACILITY",
1046
+ rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW | FACET_KEY_OPTION_EXPANDED_FILTER,
1047
+ nd_sd_journal_transform_syslog_facility,
1048
+ NULL);
1049
+
1050
+ facets_register_key_name_transformation(
1051
+ facets, "ERRNO", rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW, nd_sd_journal_transform_errno, NULL);
1052
+
1053
+ facets_register_key_name(facets, JOURNAL_KEY_ND_JOURNAL_FILE, FACET_KEY_OPTION_NEVER_FACET);
1054
+
1055
+ facets_register_key_name(facets, "SYSLOG_IDENTIFIER", rq->default_facet);
1056
+
1057
+ facets_register_key_name(facets, "UNIT", rq->default_facet);
1058
+
1059
+ facets_register_key_name(facets, "USER_UNIT", rq->default_facet);
1060
+
1061
+ facets_register_key_name_transformation(
1062
+ facets,
1063
+ "MESSAGE_ID",
1064
+ rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW | FACET_KEY_OPTION_EXPANDED_FILTER,
1065
+ nd_sd_journal_transform_message_id,
1066
+ NULL);
1067
+
1068
+ facets_register_key_name_transformation(
1069
+ facets, "_BOOT_ID", rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW, nd_sd_journal_transform_boot_id, NULL);
1070
+
1071
+ facets_register_key_name_transformation(
1072
+ facets,
1073
+ "_SYSTEMD_OWNER_UID",
1074
+ rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1075
+ nd_sd_journal_transform_uid,
1076
+ NULL);
1077
+
1078
+ facets_register_key_name_transformation(
1079
+ facets, "_UID", rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW, nd_sd_journal_transform_uid, NULL);
1080
+
1081
+ facets_register_key_name_transformation(
1082
+ facets,
1083
+ "OBJECT_SYSTEMD_OWNER_UID",
1084
+ rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW,
1085
+ nd_sd_journal_transform_uid,
1086
+ NULL);
1087
+
1088
+ facets_register_key_name_transformation(
1089
+ facets, "OBJECT_UID", rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW, nd_sd_journal_transform_uid, NULL);
1090
+
1091
+ facets_register_key_name_transformation(
1092
+ facets, "_GID", rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW, nd_sd_journal_transform_gid, NULL);
1093
+
1094
+ facets_register_key_name_transformation(
1095
+ facets, "OBJECT_GID", rq->default_facet | FACET_KEY_OPTION_TRANSFORM_VIEW, nd_sd_journal_transform_gid, NULL);
1096
+
1097
+ facets_register_key_name_transformation(
1098
+ facets, "_CAP_EFFECTIVE", FACET_KEY_OPTION_TRANSFORM_VIEW, nd_sd_journal_transform_cap_effective, NULL);
1099
+
1100
+ facets_register_key_name_transformation(
1101
+ facets, "_AUDIT_LOGINUID", FACET_KEY_OPTION_TRANSFORM_VIEW, nd_sd_journal_transform_uid, NULL);
1102
+
1103
+ facets_register_key_name_transformation(
1104
+ facets, "OBJECT_AUDIT_LOGINUID", FACET_KEY_OPTION_TRANSFORM_VIEW, nd_sd_journal_transform_uid, NULL);
1105
+
1106
+ facets_register_key_name_transformation(
1107
+ facets,
1108
+ "_SOURCE_REALTIME_TIMESTAMP",
1109
+ FACET_KEY_OPTION_TRANSFORM_VIEW,
1110
+ nd_sd_journal_transform_timestamp_usec,
1111
+ NULL);
1112
+}
1113
+
1114
+void function_systemd_journal(
1115
+ const char *transaction,
1116
+ char *function,
1117
+ usec_t *stop_monotonic_ut,
1118
+ bool *cancelled,
1119
+ BUFFER *payload,
1120
+ HTTP_ACCESS access __maybe_unused,
1121
+ const char *source __maybe_unused,
1122
+ void *data __maybe_unused)
1123
+{
1124
+ fstat_thread_calls = 0;
1125
+ fstat_thread_cached_responses = 0;
1126
+
1127
+#ifdef HAVE_SD_JOURNAL_RESTART_FIELDS
1128
+ bool have_slice = true;
1129
+#else
1130
+ bool have_slice = false;
1131
+#endif // HAVE_SD_JOURNAL_RESTART_FIELDS
1132
+
1133
+ LOGS_QUERY_STATUS tmp_fqs = {
1134
+ .facets = lqs_facets_create(
1135
+ LQS_DEFAULT_ITEMS_PER_QUERY,
1136
+ FACETS_OPTION_ALL_KEYS_FTS | FACETS_OPTION_HASH_IDS,
1137
+ SYSTEMD_ALWAYS_VISIBLE_KEYS,
1138
+ SYSTEMD_KEYS_INCLUDED_IN_FACETS,
1139
+ SYSTEMD_KEYS_EXCLUDED_FROM_FACETS,
1140
+ have_slice),
1141
+
1142
+ .rq = LOGS_QUERY_REQUEST_DEFAULTS(transaction, LQS_DEFAULT_SLICE_MODE, JOURNAL_DEFAULT_DIRECTION),
1143
+
1144
+ .cancelled = cancelled,
1145
+ .stop_monotonic_ut = stop_monotonic_ut,
1146
+ };
1147
+ LOGS_QUERY_STATUS *lqs = &tmp_fqs;
1148
+
1149
+ CLEAN_BUFFER *wb = lqs_create_output_buffer();
1150
+
1151
+ // ------------------------------------------------------------------------
1152
+ // parse the parameters
1153
+
1154
+ if (lqs_request_parse_and_validate(lqs, wb, function, payload, have_slice, "PRIORITY")) {
1155
+ systemd_journal_register_transformations(lqs);
1156
+
1157
+ // ------------------------------------------------------------------------
1158
+ // add versions to the response
1159
+
1160
+ buffer_json_journal_versions(wb);
1161
+
1162
+ // ------------------------------------------------------------------------
1163
+ // run the request
1164
+
1165
+ if (lqs->rq.info)
1166
+ lqs_info_response(wb, lqs->facets);
1167
+ else {
1168
+ nd_sd_journal_query(wb, lqs);
1169
+ if (wb->response_code == HTTP_RESP_OK)
1170
+ buffer_json_finalize(wb);
1171
+ }
1172
+ }
1173
+
1174
+ netdata_mutex_lock(&stdout_mutex);
1175
+ pluginsd_function_result_to_stdout(transaction, wb);
1176
+ netdata_mutex_unlock(&stdout_mutex);
1177
+
1178
+ lqs_cleanup(lqs);
1179
+}