master
h 867 lines 36.7 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 #ifndef NETDATA_LOGS_QUERY_STATUS_H
4 #define NETDATA_LOGS_QUERY_STATUS_H
5
6 #include "../libnetdata.h"
7
8 #define LQS_PARAMETER_HELP "help"
9 #define LQS_PARAMETER_AFTER "after"
10 #define LQS_PARAMETER_BEFORE "before"
11 #define LQS_PARAMETER_ANCHOR "anchor"
12 #define LQS_PARAMETER_LAST "last"
13 #define LQS_PARAMETER_QUERY "query"
14 #define LQS_PARAMETER_FACETS "facets"
15 #define LQS_PARAMETER_HISTOGRAM "histogram"
16 #define LQS_PARAMETER_DIRECTION "direction"
17 #define LQS_PARAMETER_IF_MODIFIED_SINCE "if_modified_since"
18 #define LQS_PARAMETER_DATA_ONLY "data_only"
19 #define LQS_PARAMETER_SOURCE "__logs_sources" // this must never conflict with user fields
20 #define LQS_PARAMETER_INFO "info"
21 #define LQS_PARAMETER_SLICE "slice"
22 #define LQS_PARAMETER_DELTA "delta"
23 #define LQS_PARAMETER_TAIL "tail"
24 #define LQS_PARAMETER_SAMPLING "sampling"
25
26 #define LQS_MAX_PARAMS 1000
27 #define LQS_DEFAULT_QUERY_DURATION (1 * 3600)
28
29 #undef LQS_SLICE_PARAMETER
30 #if LQS_DEFAULT_SLICE_MODE == 1
31 #define LQS_SLICE_PARAMETER 1
32 #endif
33
34 typedef struct {
35 const char *transaction;
36
37 FACET_KEY_OPTIONS default_facet; // the option to be used for internal fields.
38 // when the requests set facets, we disable all default facets,
39 // so that the UI has full control over them.
40
41 bool fields_are_ids; // POST works with field names, GET works with field hashes (IDs)
42 bool info; // the request is an INFO request, do not execute a query.
43
44 bool data_only; // return as fast as possible, with the requested amount of data,
45 // without scanning the entire duration.
46
47 bool slice; // apply native backend filters to slice the events database.
48 bool delta; // return incremental data for the histogram (used with data_only)
49 bool tail; // return NOT MODIFIED if no more data are available after the anchor given.
50
51 time_t after_s; // the starting timestamp of the query
52 time_t before_s; // the ending timestamp of the query
53 usec_t after_ut; // in microseconds
54 usec_t before_ut; // in microseconds
55
56 usec_t anchor; // the anchor to seek to
57 FACETS_ANCHOR_DIRECTION direction; // the direction based on the anchor (or the query timeframe)
58
59 usec_t if_modified_since; // the timestamp to check with tail == true
60
61 size_t entries; // the number of log events to return in a single response
62
63 const char *query; // full text search query string
64 const char *histogram; // the field to use for the histogram
65
66 SIMPLE_PATTERN *sources; // custom log sources to query
67 LQS_SOURCE_TYPE source_type; // pre-defined log sources to query
68
69 size_t filters; // the number of filters (facets selected) in the query
70 size_t sampling; // the number of log events to sample, when the query is too big
71
72 time_t now_s; // the timestamp the query was received
73 time_t expires_s; // the timestamp the response expires
74 } LOGS_QUERY_REQUEST;
75
76 #define LOGS_QUERY_REQUEST_DEFAULTS(function_transaction, default_slice, default_direction) \
77 (LOGS_QUERY_REQUEST) { \
78 .transaction = (function_transaction), \
79 .default_facet = FACET_KEY_OPTION_FACET, \
80 .info = false, \
81 .data_only = false, \
82 .slice = (default_slice), \
83 .delta = false, \
84 .tail = false, \
85 .after_s = 0, \
86 .before_s = 0, \
87 .anchor = 0, \
88 .if_modified_since = 0, \
89 .entries = 0, \
90 .direction = (default_direction), \
91 .query = NULL, \
92 .histogram = NULL, \
93 .sources = NULL, \
94 .source_type = LQS_SOURCE_TYPE_ALL, \
95 .filters = 0, \
96 .sampling = LQS_DEFAULT_ITEMS_SAMPLING, \
97 }
98
99 typedef struct {
100 FACETS *facets;
101
102 LOGS_QUERY_REQUEST rq;
103
104 bool *cancelled; // a pointer to the cancelling boolean
105 usec_t *stop_monotonic_ut;
106
107 struct {
108 usec_t start_ut;
109 usec_t stop_ut;
110 usec_t delta_ut;
111 } anchor;
112
113 struct {
114 usec_t start_ut;
115 usec_t stop_ut;
116 bool stop_when_full;
117 } query;
118
119 usec_t last_modified;
120
121 struct lqs_extension c;
122 } LOGS_QUERY_STATUS;
123
124 struct logs_query_data {
125 const char *transaction;
126 FACETS *facets;
127 LOGS_QUERY_REQUEST *rq;
128 BUFFER *wb;
129 };
130
131 static inline FACETS_ANCHOR_DIRECTION lgs_get_direction(const char *value) {
132 return strcasecmp(value, "forward") == 0 ? FACETS_ANCHOR_DIRECTION_FORWARD : FACETS_ANCHOR_DIRECTION_BACKWARD;
133 }
134
135 static inline void lqs_log_error(LOGS_QUERY_STATUS *lqs, const char *msg) {
136 nd_log(NDLS_COLLECTORS, NDLP_ERR,
137 "LOGS QUERY ERROR: %s, on query "
138 "timeframe [%"PRIu64" - %"PRIu64"], "
139 "anchor [%"PRIu64" - %"PRIu64"], "
140 "if_modified_since %"PRIu64", "
141 "data_only:%s, delta:%s, tail:%s, direction:%s"
142 , msg
143 , lqs->rq.after_ut
144 , lqs->rq.before_ut
145 , lqs->anchor.start_ut
146 , lqs->anchor.stop_ut
147 , lqs->rq.if_modified_since
148 , lqs->rq.data_only ? "true" : "false"
149 , lqs->rq.delta ? "true" : "false"
150 , lqs->rq.tail ? "tail" : "false"
151 , lqs->rq.direction == FACETS_ANCHOR_DIRECTION_FORWARD ? "forward" : "backward");
152 }
153
154 static inline void lqs_query_timeframe(LOGS_QUERY_STATUS *lqs, usec_t anchor_delta_ut) {
155 lqs->anchor.delta_ut = anchor_delta_ut;
156
157 if(lqs->rq.direction == FACETS_ANCHOR_DIRECTION_FORWARD) {
158 lqs->query.start_ut = (lqs->rq.data_only && lqs->anchor.start_ut) ? lqs->anchor.start_ut : lqs->rq.after_ut;
159 lqs->query.stop_ut = ((lqs->rq.data_only && lqs->anchor.stop_ut) ? lqs->anchor.stop_ut : lqs->rq.before_ut) + lqs->anchor.delta_ut;
160 }
161 else {
162 lqs->query.start_ut = ((lqs->rq.data_only && lqs->anchor.start_ut) ? lqs->anchor.start_ut : lqs->rq.before_ut) + lqs->anchor.delta_ut;
163 lqs->query.stop_ut = (lqs->rq.data_only && lqs->anchor.stop_ut) ? lqs->anchor.stop_ut : lqs->rq.after_ut;
164 }
165
166 lqs->query.stop_when_full = (lqs->rq.data_only && !lqs->anchor.stop_ut);
167 }
168
169 static inline void lqs_function_help(LOGS_QUERY_STATUS *lqs, BUFFER *wb) {
170 buffer_reset(wb);
171 wb->content_type = CT_TEXT_PLAIN;
172 wb->response_code = HTTP_RESP_OK;
173
174 buffer_sprintf(wb,
175 "%s / %s\n"
176 "\n"
177 "%s\n"
178 "\n"
179 "The following parameters are supported:\n"
180 "\n"
181 , program_name
182 , LQS_FUNCTION_NAME
183 , LQS_FUNCTION_DESCRIPTION
184 );
185
186 buffer_sprintf(wb,
187 " " LQS_PARAMETER_HELP "\n"
188 " Shows this help message.\n"
189 "\n"
190 );
191
192 buffer_sprintf(wb,
193 " " LQS_PARAMETER_INFO "\n"
194 " Request initial configuration information about the plugin.\n"
195 " The key entity returned is the required_params array, which includes\n"
196 " all the available log sources.\n"
197 " When `" LQS_PARAMETER_INFO "` is requested, all other parameters are ignored.\n"
198 "\n"
199 );
200
201 buffer_sprintf(wb,
202 " " LQS_PARAMETER_DATA_ONLY ":true or " LQS_PARAMETER_DATA_ONLY ":false\n"
203 " Quickly respond with data requested, without generating a\n"
204 " `histogram`, `facets` counters and `items`.\n"
205 "\n"
206 );
207
208 buffer_sprintf(wb,
209 " " LQS_PARAMETER_DELTA ":true or " LQS_PARAMETER_DELTA ":false\n"
210 " When doing data only queries, include deltas for histogram, facets and items.\n"
211 "\n"
212 );
213
214 buffer_sprintf(wb,
215 " " LQS_PARAMETER_TAIL ":true or " LQS_PARAMETER_TAIL ":false\n"
216 " When doing data only queries, respond with the newest messages,\n"
217 " and up to the anchor, but calculate deltas (if requested) for\n"
218 " the duration [anchor - before].\n"
219 "\n"
220 );
221
222 #ifdef LQS_SLICE_PARAMETER
223 buffer_sprintf(wb,
224 " " LQS_PARAMETER_SLICE ":true or " LQS_PARAMETER_SLICE ":false\n"
225 " When it is turned on, the plugin is is slicing the logs database,\n"
226 " utilizing the underlying available indexes.\n"
227 " When it is off, all filtering is done by the plugin.\n"
228 " The default is: %s\n"
229 "\n"
230 , lqs->rq.slice ? "true" : "false"
231 );
232 #endif
233 buffer_sprintf(wb,
234 " " LQS_PARAMETER_SOURCE ":SOURCE\n"
235 " Query only the specified log sources.\n"
236 " Do an `" LQS_PARAMETER_INFO "` query to find the sources.\n"
237 "\n"
238 );
239
240 buffer_sprintf(wb,
241 " " LQS_PARAMETER_BEFORE ":TIMESTAMP_IN_SECONDS\n"
242 " Absolute or relative (to now) timestamp in seconds, to start the query.\n"
243 " The query is always executed from the most recent to the oldest log entry.\n"
244 " If not given the default is: now.\n"
245 "\n"
246 );
247
248 buffer_sprintf(wb,
249 " " LQS_PARAMETER_AFTER ":TIMESTAMP_IN_SECONDS\n"
250 " Absolute or relative (to `before`) timestamp in seconds, to end the query.\n"
251 " If not given, the default is %d.\n"
252 "\n"
253 , -LQS_DEFAULT_QUERY_DURATION
254 );
255
256 buffer_sprintf(wb,
257 " " LQS_PARAMETER_LAST ":ITEMS\n"
258 " The number of items to return.\n"
259 " The default is %zu.\n"
260 "\n"
261 , lqs->rq.entries
262 );
263
264 buffer_sprintf(wb,
265 " " LQS_PARAMETER_SAMPLING ":ITEMS\n"
266 " The number of log entries to sample to estimate facets counters and histogram.\n"
267 " The default is %zu.\n"
268 "\n"
269 , lqs->rq.sampling
270 );
271
272 buffer_sprintf(wb,
273 " " LQS_PARAMETER_ANCHOR ":TIMESTAMP_IN_MICROSECONDS\n"
274 " Return items relative to this timestamp.\n"
275 " The exact items to be returned depend on the query `" LQS_PARAMETER_DIRECTION "`.\n"
276 "\n"
277 );
278
279 buffer_sprintf(wb,
280 " " LQS_PARAMETER_DIRECTION ":forward or " LQS_PARAMETER_DIRECTION ":backward\n"
281 " When set to `backward` (default) the items returned are the newest before the\n"
282 " `" LQS_PARAMETER_ANCHOR "`, (or `" LQS_PARAMETER_BEFORE "` if `" LQS_PARAMETER_ANCHOR "` is not set)\n"
283 " When set to `forward` the items returned are the oldest after the\n"
284 " `" LQS_PARAMETER_ANCHOR "`, (or `" LQS_PARAMETER_AFTER "` if `" LQS_PARAMETER_ANCHOR "` is not set)\n"
285 " The default is: %s\n"
286 "\n"
287 , lqs->rq.direction == FACETS_ANCHOR_DIRECTION_FORWARD ? "forward" : "backward"
288 );
289
290 buffer_sprintf(wb,
291 " " LQS_PARAMETER_QUERY ":SIMPLE_PATTERN\n"
292 " Do a full text search to find the log entries matching the pattern given.\n"
293 " The plugin is searching for matches on all fields of the database.\n"
294 "\n"
295 );
296
297 buffer_sprintf(wb,
298 " " LQS_PARAMETER_IF_MODIFIED_SINCE ":TIMESTAMP_IN_MICROSECONDS\n"
299 " Each successful response, includes a `last_modified` field.\n"
300 " By providing the timestamp to the `" LQS_PARAMETER_IF_MODIFIED_SINCE "` parameter,\n"
301 " the plugin will return 200 with a successful response, or 304 if the source has not\n"
302 " been modified since that timestamp.\n"
303 "\n"
304 );
305
306 buffer_sprintf(wb,
307 " " LQS_PARAMETER_HISTOGRAM ":facet_id\n"
308 " Use the given `facet_id` for the histogram.\n"
309 " This parameter is ignored in `" LQS_PARAMETER_DATA_ONLY "` mode.\n"
310 "\n"
311 );
312
313 buffer_sprintf(wb,
314 " " LQS_PARAMETER_FACETS ":facet_id1,facet_id2,facet_id3,...\n"
315 " Add the given facets to the list of fields for which analysis is required.\n"
316 " The plugin will offer both a histogram and facet value counters for its values.\n"
317 " This parameter is ignored in `" LQS_PARAMETER_DATA_ONLY "` mode.\n"
318 "\n"
319 );
320
321 buffer_sprintf(wb,
322 " facet_id:value_id1,value_id2,value_id3,...\n"
323 " Apply filters to the query, based on the facet IDs returned.\n"
324 " Each `facet_id` can be given once, but multiple `facet_ids` can be given.\n"
325 "\n"
326 );
327 }
328
329 static inline bool lqs_request_parse_json_payload(json_object *jobj, void *data, BUFFER *error) {
330 const char *path = "";
331 struct logs_query_data *qd = data;
332 LOGS_QUERY_REQUEST *rq = qd->rq;
333 BUFFER *wb = qd->wb;
334 FACETS *facets = qd->facets;
335 // const char *transaction = qd->transaction;
336
337 buffer_flush(error);
338
339 JSONC_PARSE_BOOL_OR_ERROR_AND_RETURN(jobj, path, LQS_PARAMETER_INFO, rq->info, error, JSONC_OPTIONAL);
340 JSONC_PARSE_BOOL_OR_ERROR_AND_RETURN(jobj, path, LQS_PARAMETER_DELTA, rq->delta, error, JSONC_OPTIONAL);
341 JSONC_PARSE_BOOL_OR_ERROR_AND_RETURN(jobj, path, LQS_PARAMETER_TAIL, rq->tail, error, JSONC_OPTIONAL);
342 JSONC_PARSE_BOOL_OR_ERROR_AND_RETURN(jobj, path, LQS_PARAMETER_SLICE, rq->slice, error, JSONC_OPTIONAL);
343 JSONC_PARSE_BOOL_OR_ERROR_AND_RETURN(jobj, path, LQS_PARAMETER_DATA_ONLY, rq->data_only, error, JSONC_OPTIONAL);
344 JSONC_PARSE_UINT64_OR_ERROR_AND_RETURN(jobj, path, LQS_PARAMETER_SAMPLING, rq->sampling, error, JSONC_OPTIONAL);
345 JSONC_PARSE_INT64_OR_ERROR_AND_RETURN(jobj, path, LQS_PARAMETER_AFTER, rq->after_s, error, JSONC_OPTIONAL);
346 JSONC_PARSE_INT64_OR_ERROR_AND_RETURN(jobj, path, LQS_PARAMETER_BEFORE, rq->before_s, error, JSONC_OPTIONAL);
347 JSONC_PARSE_UINT64_OR_ERROR_AND_RETURN(jobj, path, LQS_PARAMETER_IF_MODIFIED_SINCE, rq->if_modified_since, error, JSONC_OPTIONAL);
348 JSONC_PARSE_UINT64_OR_ERROR_AND_RETURN(jobj, path, LQS_PARAMETER_ANCHOR, rq->anchor, error, JSONC_OPTIONAL);
349 JSONC_PARSE_UINT64_OR_ERROR_AND_RETURN(jobj, path, LQS_PARAMETER_LAST, rq->entries, error, JSONC_OPTIONAL);
350 JSONC_PARSE_TXT2ENUM_OR_ERROR_AND_RETURN(jobj, path, LQS_PARAMETER_DIRECTION, lgs_get_direction, rq->direction, error, JSONC_OPTIONAL);
351 JSONC_PARSE_TXT2STRDUPZ_OR_ERROR_AND_RETURN(jobj, path, LQS_PARAMETER_QUERY, rq->query, error, JSONC_OPTIONAL);
352 JSONC_PARSE_TXT2STRDUPZ_OR_ERROR_AND_RETURN(jobj, path, LQS_PARAMETER_HISTOGRAM, rq->histogram, error, JSONC_OPTIONAL);
353
354 json_object *fcts;
355 if (json_object_object_get_ex(jobj, LQS_PARAMETER_FACETS, &fcts)) {
356 if (json_object_get_type(fcts) != json_type_array) {
357 buffer_sprintf(error, "member '%s' is not an array.", LQS_PARAMETER_FACETS);
358 // nd_log(NDLS_COLLECTORS, NDLP_ERR, "POST payload: '%s' is not an array", LQS_PARAMETER_FACETS);
359 return false;
360 }
361
362 rq->default_facet = FACET_KEY_OPTION_NONE;
363 facets_reset_and_disable_all_facets(facets);
364
365 buffer_json_member_add_array(wb, LQS_PARAMETER_FACETS);
366
367 size_t facets_len = json_object_array_length(fcts);
368 for (size_t i = 0; i < facets_len; i++) {
369 json_object *fct = json_object_array_get_idx(fcts, i);
370
371 if (json_object_get_type(fct) != json_type_string) {
372 buffer_sprintf(error, "facets array item %zu is not a string", i);
373 // nd_log(NDLS_COLLECTORS, NDLP_ERR, "POST payload: facets array item %zu is not a string", i);
374 return false;
375 }
376
377 const char *value = json_object_get_string(fct);
378 facets_register_facet(facets, value, FACET_KEY_OPTION_FACET|FACET_KEY_OPTION_FTS|FACET_KEY_OPTION_REORDER);
379 buffer_json_add_array_item_string(wb, value);
380 }
381
382 buffer_json_array_close(wb); // facets
383 }
384
385 json_object *selections;
386 if (json_object_object_get_ex(jobj, "selections", &selections)) {
387 if (json_object_get_type(selections) != json_type_object) {
388 buffer_sprintf(error, "member 'selections' is not an object");
389 // nd_log(NDLS_COLLECTORS, NDLP_ERR, "POST payload: '%s' is not an object", "selections");
390 return false;
391 }
392
393 buffer_json_member_add_object(wb, "selections");
394
395 CLEAN_BUFFER *sources_list = buffer_create(0, NULL);
396
397 json_object_object_foreach(selections, key, val) {
398 if(strcmp(key, "query") == 0) continue;
399
400 if (json_object_get_type(val) != json_type_array) {
401 buffer_sprintf(error, "selection '%s' is not an array", key);
402 // nd_log(NDLS_COLLECTORS, NDLP_ERR, "POST payload: selection '%s' is not an array", key);
403 return false;
404 }
405
406 bool is_source = false;
407 if(strcmp(key, LQS_PARAMETER_SOURCE) == 0) {
408 // reset the sources, so that only what the user selects will be shown
409 is_source = true;
410 rq->source_type = LQS_SOURCE_TYPE_NONE;
411 }
412
413 buffer_json_member_add_array(wb, key);
414
415 size_t values_len = json_object_array_length(val);
416 for (size_t i = 0; i < values_len; i++) {
417 json_object *value_obj = json_object_array_get_idx(val, i);
418
419 if (json_object_get_type(value_obj) != json_type_string) {
420 buffer_sprintf(error, "selection '%s' array item %zu is not a string", key, i);
421 // nd_log(NDLS_COLLECTORS, NDLP_ERR, "POST payload: selection '%s' array item %zu is not a string", key, i);
422 return false;
423 }
424
425 const char *value = json_object_get_string(value_obj);
426
427 if(is_source) {
428 // processing sources
429 LQS_SOURCE_TYPE t = LQS_FUNCTION_GET_INTERNAL_SOURCE_TYPE(value);
430 if(t != LQS_SOURCE_TYPE_NONE) {
431 rq->source_type |= t;
432 value = NULL;
433 }
434 else {
435 // else, match the source, whatever it is
436 if(buffer_strlen(sources_list))
437 buffer_putc(sources_list, '|');
438
439 buffer_strcat(sources_list, value);
440 }
441 }
442 else {
443 // Call facets_register_facet_id_filter for each value
444 facets_register_facet_filter(
445 facets, key, value, FACET_KEY_OPTION_FTS | FACET_KEY_OPTION_REORDER);
446
447 rq->filters++;
448 }
449
450 buffer_json_add_array_item_string(wb, value);
451 }
452
453 buffer_json_array_close(wb); // key
454 }
455
456 if(buffer_strlen(sources_list)) {
457 simple_pattern_free(rq->sources);
458 rq->sources = simple_pattern_create(buffer_tostring(sources_list), "|", SIMPLE_PATTERN_EXACT, false);
459 }
460
461 buffer_json_object_close(wb); // selections
462 }
463
464 facets_use_hashes_for_ids(facets, false);
465 rq->fields_are_ids = false;
466 return true;
467 }
468
469 static inline bool lqs_request_parse_POST(LOGS_QUERY_STATUS *lqs, BUFFER *wb, BUFFER *payload, const char *transaction) {
470 FACETS *facets = lqs->facets;
471 LOGS_QUERY_REQUEST *rq = &lqs->rq;
472
473 buffer_json_member_add_object(wb, "_request");
474
475 struct logs_query_data qd = {
476 .transaction = transaction,
477 .facets = facets,
478 .rq = rq,
479 .wb = wb,
480 };
481
482 int code;
483 CLEAN_JSON_OBJECT *jobj =
484 json_parse_function_payload_or_error(wb, payload, &code, lqs_request_parse_json_payload, &qd);
485 wb->response_code = code;
486
487 return (jobj && code == HTTP_RESP_OK);
488 }
489
490 static inline bool lqs_request_parse_GET(LOGS_QUERY_STATUS *lqs, BUFFER *wb, char *function) {
491 FACETS *facets = lqs->facets;
492 LOGS_QUERY_REQUEST *rq = &lqs->rq;
493
494 buffer_json_member_add_object(wb, "_request");
495
496 CLEAN_CHAR_P *func_copy = strdupz(function);
497
498 char *words[LQS_MAX_PARAMS] = { NULL };
499 size_t num_words = quoted_strings_splitter_whitespace(func_copy, words, LQS_MAX_PARAMS);
500 for(int i = 1; i < LQS_MAX_PARAMS;i++) {
501 char *keyword = get_word(words, num_words, i);
502 if(!keyword) break;
503
504 if(strcmp(keyword, LQS_PARAMETER_HELP) == 0) {
505 lqs_function_help(lqs, wb);
506 return false;
507 }
508 else if(strcmp(keyword, LQS_PARAMETER_INFO) == 0) {
509 rq->info = true;
510 }
511 else if(strncmp(keyword, LQS_PARAMETER_DELTA ":", sizeof(LQS_PARAMETER_DELTA ":") - 1) == 0) {
512 char *v = &keyword[sizeof(LQS_PARAMETER_DELTA ":") - 1];
513
514 if(strcmp(v, "false") == 0 || strcmp(v, "no") == 0 || strcmp(v, "0") == 0)
515 rq->delta = false;
516 else
517 rq->delta = true;
518 }
519 else if(strncmp(keyword, LQS_PARAMETER_TAIL ":", sizeof(LQS_PARAMETER_TAIL ":") - 1) == 0) {
520 char *v = &keyword[sizeof(LQS_PARAMETER_TAIL ":") - 1];
521
522 if(strcmp(v, "false") == 0 || strcmp(v, "no") == 0 || strcmp(v, "0") == 0)
523 rq->tail = false;
524 else
525 rq->tail = true;
526 }
527 else if(strncmp(keyword, LQS_PARAMETER_SAMPLING ":", sizeof(LQS_PARAMETER_SAMPLING ":") - 1) == 0) {
528 rq->sampling = str2ul(&keyword[sizeof(LQS_PARAMETER_SAMPLING ":") - 1]);
529 }
530 else if(strncmp(keyword, LQS_PARAMETER_DATA_ONLY ":", sizeof(LQS_PARAMETER_DATA_ONLY ":") - 1) == 0) {
531 char *v = &keyword[sizeof(LQS_PARAMETER_DATA_ONLY ":") - 1];
532
533 if(strcmp(v, "false") == 0 || strcmp(v, "no") == 0 || strcmp(v, "0") == 0)
534 rq->data_only = false;
535 else
536 rq->data_only = true;
537 }
538 else if(strncmp(keyword, LQS_PARAMETER_SLICE ":", sizeof(LQS_PARAMETER_SLICE ":") - 1) == 0) {
539 char *v = &keyword[sizeof(LQS_PARAMETER_SLICE ":") - 1];
540
541 if(strcmp(v, "false") == 0 || strcmp(v, "no") == 0 || strcmp(v, "0") == 0)
542 rq->slice = false;
543 else
544 rq->slice = true;
545 }
546 else if(strncmp(keyword, LQS_PARAMETER_SOURCE ":", sizeof(LQS_PARAMETER_SOURCE ":") - 1) == 0) {
547 char *value = &keyword[sizeof(LQS_PARAMETER_SOURCE ":") - 1];
548
549 buffer_json_member_add_array(wb, LQS_PARAMETER_SOURCE);
550
551 CLEAN_BUFFER *sources_list = buffer_create(0, NULL);
552
553 rq->source_type = LQS_SOURCE_TYPE_NONE;
554 while(value) {
555 char *sep = strchr(value, ',');
556 if(sep)
557 *sep++ = '\0';
558
559 buffer_json_add_array_item_string(wb, value);
560
561 LQS_SOURCE_TYPE t = LQS_FUNCTION_GET_INTERNAL_SOURCE_TYPE(value);
562 if(t != LQS_SOURCE_TYPE_NONE) {
563 rq->source_type |= t;
564 }
565 else {
566 // else, match the source, whatever it is
567 if(buffer_strlen(sources_list))
568 buffer_putc(sources_list, '|');
569
570 buffer_strcat(sources_list, value);
571 }
572
573 value = sep;
574 }
575
576 if(buffer_strlen(sources_list)) {
577 simple_pattern_free(rq->sources);
578 rq->sources = simple_pattern_create(buffer_tostring(sources_list), "|", SIMPLE_PATTERN_EXACT, false);
579 }
580
581 buffer_json_array_close(wb); // source
582 }
583 else if(strncmp(keyword, LQS_PARAMETER_AFTER ":", sizeof(LQS_PARAMETER_AFTER ":") - 1) == 0) {
584 rq->after_s = str2l(&keyword[sizeof(LQS_PARAMETER_AFTER ":") - 1]);
585 }
586 else if(strncmp(keyword, LQS_PARAMETER_BEFORE ":", sizeof(LQS_PARAMETER_BEFORE ":") - 1) == 0) {
587 rq->before_s = str2l(&keyword[sizeof(LQS_PARAMETER_BEFORE ":") - 1]);
588 }
589 else if(strncmp(keyword, LQS_PARAMETER_IF_MODIFIED_SINCE ":", sizeof(LQS_PARAMETER_IF_MODIFIED_SINCE ":") - 1) == 0) {
590 rq->if_modified_since = str2ull(&keyword[sizeof(LQS_PARAMETER_IF_MODIFIED_SINCE ":") - 1], NULL);
591 }
592 else if(strncmp(keyword, LQS_PARAMETER_ANCHOR ":", sizeof(LQS_PARAMETER_ANCHOR ":") - 1) == 0) {
593 rq->anchor = str2ull(&keyword[sizeof(LQS_PARAMETER_ANCHOR ":") - 1], NULL);
594 }
595 else if(strncmp(keyword, LQS_PARAMETER_DIRECTION ":", sizeof(LQS_PARAMETER_DIRECTION ":") - 1) == 0) {
596 rq->direction = lgs_get_direction(&keyword[sizeof(LQS_PARAMETER_DIRECTION ":") - 1]);
597 }
598 else if(strncmp(keyword, LQS_PARAMETER_LAST ":", sizeof(LQS_PARAMETER_LAST ":") - 1) == 0) {
599 rq->entries = str2ul(&keyword[sizeof(LQS_PARAMETER_LAST ":") - 1]);
600 }
601 else if(strncmp(keyword, LQS_PARAMETER_QUERY ":", sizeof(LQS_PARAMETER_QUERY ":") - 1) == 0) {
602 freez((void *)rq->query);
603 rq->query= strdupz(&keyword[sizeof(LQS_PARAMETER_QUERY ":") - 1]);
604 }
605 else if(strncmp(keyword, LQS_PARAMETER_HISTOGRAM ":", sizeof(LQS_PARAMETER_HISTOGRAM ":") - 1) == 0) {
606 freez((void *)rq->histogram);
607 rq->histogram = strdupz(&keyword[sizeof(LQS_PARAMETER_HISTOGRAM ":") - 1]);
608 }
609 else if(strncmp(keyword, LQS_PARAMETER_FACETS ":", sizeof(LQS_PARAMETER_FACETS ":") - 1) == 0) {
610 rq->default_facet = FACET_KEY_OPTION_NONE;
611 facets_reset_and_disable_all_facets(facets);
612
613 char *value = &keyword[sizeof(LQS_PARAMETER_FACETS ":") - 1];
614 if(*value) {
615 buffer_json_member_add_array(wb, LQS_PARAMETER_FACETS);
616
617 while(value) {
618 char *sep = strchr(value, ',');
619 if(sep)
620 *sep++ = '\0';
621
622 facets_register_facet_id(facets, value, FACET_KEY_OPTION_FACET|FACET_KEY_OPTION_FTS|FACET_KEY_OPTION_REORDER);
623 buffer_json_add_array_item_string(wb, value);
624
625 value = sep;
626 }
627
628 buffer_json_array_close(wb); // facets
629 }
630 }
631 else {
632 char *value = strchr(keyword, ':');
633 if(value) {
634 *value++ = '\0';
635
636 buffer_json_member_add_array(wb, keyword);
637
638 while(value) {
639 char *sep = strchr(value, ',');
640 if(sep)
641 *sep++ = '\0';
642
643 facets_register_facet_filter_id(
644 facets, keyword, value,
645 FACET_KEY_OPTION_FTS | FACET_KEY_OPTION_REORDER);
646
647 buffer_json_add_array_item_string(wb, value);
648 rq->filters++;
649
650 value = sep;
651 }
652
653 buffer_json_array_close(wb); // keyword
654 }
655 }
656 }
657
658 facets_use_hashes_for_ids(facets, true);
659 rq->fields_are_ids = true;
660 return true;
661 }
662
663 static inline void lqs_info_response(BUFFER *wb, FACETS *facets) {
664 // the buffer already has the request in it
665 // DO NOT FLUSH IT
666
667 buffer_json_member_add_uint64(wb, "v", 3);
668 facets_accepted_parameters_to_json_array(facets, wb, false);
669 buffer_json_member_add_array(wb, "required_params");
670 {
671 buffer_json_add_array_item_object(wb);
672 {
673 buffer_json_member_add_string(wb, "id", LQS_PARAMETER_SOURCE);
674 buffer_json_member_add_string(wb, "name", LQS_PARAMETER_SOURCE_NAME);
675 buffer_json_member_add_string(wb, "help", "Select the logs source to query");
676 buffer_json_member_add_string(wb, "type", "multiselect");
677 buffer_json_member_add_array(wb, "options");
678 {
679 LQS_FUNCTION_SOURCE_TO_JSON_ARRAY(wb);
680 }
681 buffer_json_array_close(wb); // options array
682 }
683 buffer_json_object_close(wb); // required params object
684 }
685 buffer_json_array_close(wb); // required_params array
686
687 facets_table_config(facets, wb);
688
689 buffer_json_member_add_uint64(wb, "status", HTTP_RESP_OK);
690 buffer_json_member_add_string(wb, "type", "table");
691 buffer_json_member_add_string(wb, "help", LQS_FUNCTION_DESCRIPTION);
692 buffer_json_finalize(wb);
693
694 wb->content_type = CT_APPLICATION_JSON;
695 wb->response_code = HTTP_RESP_OK;
696 }
697
698 static inline BUFFER *lqs_create_output_buffer(void) {
699 BUFFER *wb = buffer_create(0, NULL);
700 buffer_reset(wb);
701 buffer_json_initialize(wb, "\"", "\"", 0, true, BUFFER_JSON_OPTIONS_MINIFY);
702 return wb;
703 }
704
705 static inline FACETS *lqs_facets_create(uint32_t items_to_return, FACETS_OPTIONS options, const char *visible_keys, const char *facet_keys, const char *non_facet_keys, bool have_slice) {
706 FACETS *facets = facets_create(items_to_return, options,
707 visible_keys, facet_keys, non_facet_keys);
708
709 facets_accepted_param(facets, LQS_PARAMETER_INFO);
710 facets_accepted_param(facets, LQS_PARAMETER_SOURCE);
711 facets_accepted_param(facets, LQS_PARAMETER_AFTER);
712 facets_accepted_param(facets, LQS_PARAMETER_BEFORE);
713 facets_accepted_param(facets, LQS_PARAMETER_ANCHOR);
714 facets_accepted_param(facets, LQS_PARAMETER_DIRECTION);
715 facets_accepted_param(facets, LQS_PARAMETER_LAST);
716 facets_accepted_param(facets, LQS_PARAMETER_QUERY);
717 facets_accepted_param(facets, LQS_PARAMETER_FACETS);
718 facets_accepted_param(facets, LQS_PARAMETER_HISTOGRAM);
719 facets_accepted_param(facets, LQS_PARAMETER_IF_MODIFIED_SINCE);
720 facets_accepted_param(facets, LQS_PARAMETER_DATA_ONLY);
721 facets_accepted_param(facets, LQS_PARAMETER_DELTA);
722 facets_accepted_param(facets, LQS_PARAMETER_TAIL);
723 facets_accepted_param(facets, LQS_PARAMETER_SAMPLING);
724
725 if(have_slice)
726 facets_accepted_param(facets, LQS_PARAMETER_SLICE);
727
728 return facets;
729 }
730
731 static inline bool lqs_request_parse_and_validate(LOGS_QUERY_STATUS *lqs, BUFFER *wb, char *function, BUFFER *payload, bool have_slice, const char *default_histogram) {
732 LOGS_QUERY_REQUEST *rq = &lqs->rq;
733 FACETS *facets = lqs->facets;
734
735 if( (payload && !lqs_request_parse_POST(lqs, wb, payload, rq->transaction)) ||
736 (!payload && !lqs_request_parse_GET(lqs, wb, function)) )
737 return false;
738
739 // ----------------------------------------------------------------------------------------------------------------
740 // validate parameters
741
742 if(rq->query && !*rq->query) {
743 freez((void *)rq->query);
744 rq->query = NULL;
745 }
746
747 if(rq->histogram && !*rq->histogram) {
748 freez((void *)rq->histogram);
749 rq->histogram = NULL;
750 }
751
752 if(!rq->data_only)
753 rq->delta = false;
754
755 if(!rq->data_only || !rq->if_modified_since)
756 rq->tail = false;
757
758 rq->now_s = now_realtime_sec();
759 rq->expires_s = rq->now_s + 1;
760 wb->expires = rq->expires_s;
761
762 if(!rq->after_s && !rq->before_s) {
763 rq->before_s = rq->now_s;
764 rq->after_s = rq->before_s - LQS_DEFAULT_QUERY_DURATION;
765 }
766 else
767 rrdr_relative_window_to_absolute(&rq->after_s, &rq->before_s, rq->now_s);
768
769 if(rq->after_s > rq->before_s) {
770 time_t tmp = rq->after_s;
771 rq->after_s = rq->before_s;
772 rq->before_s = tmp;
773 }
774
775 if(rq->after_s == rq->before_s)
776 rq->after_s = rq->before_s - LQS_DEFAULT_QUERY_DURATION;
777
778 rq->after_ut = rq->after_s * USEC_PER_SEC;
779 rq->before_ut = (rq->before_s * USEC_PER_SEC) + USEC_PER_SEC - 1;
780
781 if(!rq->entries)
782 rq->entries = LQS_DEFAULT_ITEMS_PER_QUERY;
783
784 // ----------------------------------------------------------------------------------------------------------------
785 // validate the anchor
786
787 lqs->last_modified = 0;
788 lqs->anchor.start_ut = lqs->rq.anchor;
789 lqs->anchor.stop_ut = 0;
790
791 if(lqs->anchor.start_ut && lqs->rq.tail) {
792 // a tail request
793 // we need the top X entries from BEFORE
794 // but, we need to calculate the facets and the
795 // histogram up to the anchor
796 lqs->rq.direction = FACETS_ANCHOR_DIRECTION_BACKWARD;
797 lqs->anchor.start_ut = 0;
798 lqs->anchor.stop_ut = lqs->rq.anchor;
799 }
800
801 if(lqs->rq.anchor && lqs->rq.anchor < lqs->rq.after_ut) {
802 lqs_log_error(lqs, "received anchor is too small for query timeframe, ignoring anchor");
803 lqs->rq.anchor = 0;
804 lqs->anchor.start_ut = 0;
805 lqs->anchor.stop_ut = 0;
806 lqs->rq.direction = FACETS_ANCHOR_DIRECTION_BACKWARD;
807 }
808 else if(lqs->rq.anchor > lqs->rq.before_ut) {
809 lqs_log_error(lqs, "received anchor is too big for query timeframe, ignoring anchor");
810 lqs->rq.anchor = 0;
811 lqs->anchor.start_ut = 0;
812 lqs->anchor.stop_ut = 0;
813 lqs->rq.direction = FACETS_ANCHOR_DIRECTION_BACKWARD;
814 }
815
816 facets_set_anchor(facets, lqs->anchor.start_ut, lqs->anchor.stop_ut, lqs->rq.direction);
817
818 facets_set_additional_options(facets,
819 ((lqs->rq.data_only) ? FACETS_OPTION_DATA_ONLY : 0) |
820 ((lqs->rq.delta) ? FACETS_OPTION_SHOW_DELTAS : 0));
821
822 facets_set_items(facets, lqs->rq.entries);
823 facets_set_query(facets, lqs->rq.query);
824
825 if(lqs->rq.slice && have_slice)
826 facets_enable_slice_mode(facets);
827 else
828 lqs->rq.slice = false;
829
830 if(lqs->rq.histogram) {
831 if(lqs->rq.fields_are_ids)
832 facets_set_timeframe_and_histogram_by_id(facets, lqs->rq.histogram, lqs->rq.after_ut, lqs->rq.before_ut);
833 else
834 facets_set_timeframe_and_histogram_by_name(facets, lqs->rq.histogram, lqs->rq.after_ut, lqs->rq.before_ut);
835 }
836 else if(default_histogram)
837 facets_set_timeframe_and_histogram_by_name(facets, default_histogram, lqs->rq.after_ut, lqs->rq.before_ut);
838
839 // complete the request object
840 buffer_json_member_add_boolean(wb, LQS_PARAMETER_INFO, lqs->rq.info);
841 buffer_json_member_add_boolean(wb, LQS_PARAMETER_SLICE, lqs->rq.slice);
842 buffer_json_member_add_boolean(wb, LQS_PARAMETER_DATA_ONLY, lqs->rq.data_only);
843 buffer_json_member_add_boolean(wb, LQS_PARAMETER_DELTA, lqs->rq.delta);
844 buffer_json_member_add_boolean(wb, LQS_PARAMETER_TAIL, lqs->rq.tail);
845 buffer_json_member_add_uint64(wb, LQS_PARAMETER_SAMPLING, lqs->rq.sampling);
846 buffer_json_member_add_uint64(wb, "source_type", lqs->rq.source_type);
847 buffer_json_member_add_uint64(wb, LQS_PARAMETER_AFTER, lqs->rq.after_ut / USEC_PER_SEC);
848 buffer_json_member_add_uint64(wb, LQS_PARAMETER_BEFORE, lqs->rq.before_ut / USEC_PER_SEC);
849 buffer_json_member_add_uint64(wb, "if_modified_since", lqs->rq.if_modified_since);
850 buffer_json_member_add_uint64(wb, LQS_PARAMETER_ANCHOR, lqs->rq.anchor);
851 buffer_json_member_add_string(wb, LQS_PARAMETER_DIRECTION, lqs->rq.direction == FACETS_ANCHOR_DIRECTION_FORWARD ? "forward" : "backward");
852 buffer_json_member_add_uint64(wb, LQS_PARAMETER_LAST, lqs->rq.entries);
853 buffer_json_member_add_string(wb, LQS_PARAMETER_QUERY, lqs->rq.query);
854 buffer_json_member_add_string(wb, LQS_PARAMETER_HISTOGRAM, lqs->rq.histogram);
855 buffer_json_object_close(wb); // request
856
857 return true;
858 }
859
860 static inline void lqs_cleanup(LOGS_QUERY_STATUS *lqs) {
861 freez((void *)lqs->rq.query);
862 freez((void *)lqs->rq.histogram);
863 simple_pattern_free(lqs->rq.sources);
864 facets_destroy(lqs->facets);
865 }
866
867 #endif //NETDATA_LOGS_QUERY_STATUS_H