remove the line length limit from pluginsd (#16013)
* remove the line length limit from pluginsd * initialize the buffer on every iteration * buffer_tostring inlined * Release buffer --------- Co-authored-by: Stelios Fragkakis <52996999+stelfrag@users.noreply.github.com>
Costa Tsaousis committed
Sep 21, 2023 at 19:05 UTC
060f3ff7bf5b91899280ceaefaddf7ff4016a1f6
7 files changed
+47
-45
collectors/plugins.d/pluginsd_parser.c
+10
-3
@@ -2352,15 +2352,22 @@ inline size_t pluginsd_process(RRDHOST *host, struct plugind *cd, FILE *fp_plugi
2352
netdata_thread_cleanup_push(pluginsd_process_thread_cleanup, parser);
2353
2354
buffered_reader_init(&parser->reader);
2355
- char buffer[PLUGINSD_LINE_MAX + 2];
2355
+ BUFFER *buffer = buffer_create(sizeof(parser->reader.read_buffer) + 2, NULL);
2356
while(likely(service_running(SERVICE_COLLECTORS))) {
2357
- if (unlikely(!buffered_reader_next_line(&parser->reader, buffer, PLUGINSD_LINE_MAX + 2))) {
2357
+ if (unlikely(!buffered_reader_next_line(&parser->reader, buffer))) {
2358
if(unlikely(!buffered_reader_read_timeout(&parser->reader, fileno((FILE *)parser->fp_input), 2 * 60 * MSEC_PER_SEC)))
2359
break;
2360
+
2361
+ continue;
2362
}
2361
- else if(unlikely(parser_action(parser, buffer)))
2363
+
2364
+ if(unlikely(parser_action(parser, buffer->buffer)))
2365
break;
2366
+
2367
+ buffer->len = 0;
2368
+ buffer->buffer[0] = '\0';
2369
}
2370
+ buffer_free(buffer);
2371
2372
cd->unsafe.enabled = parser->user.enabled;
2373
count = parser->user.data_collections_count;
collectors/plugins.d/pluginsd_parser.h
+3
-3
@@ -154,8 +154,8 @@ static inline int parser_action(PARSER *parser, char *input) {
154
parser->line++;
155
156
if(unlikely(parser->flags & PARSER_DEFER_UNTIL_KEYWORD)) {
157
- char command[PLUGINSD_LINE_MAX + 1];
158
- bool has_keyword = find_first_keyword(input, command, PLUGINSD_LINE_MAX, isspace_map_pluginsd);
157
+ char command[100 + 1];
158
+ bool has_keyword = find_first_keyword(input, command, 100, isspace_map_pluginsd);
159
160
if(!has_keyword || strcmp(command, parser->defer.end_keyword) != 0) {
161
if(parser->defer.response) {
@@ -183,7 +183,7 @@ static inline int parser_action(PARSER *parser, char *input) {
183
return 0;
184
}
185
186
- char *words[PLUGINSD_MAX_WORDS];
186
+ static __thread char *words[PLUGINSD_MAX_WORDS];
187
size_t num_words = quoted_strings_splitter_pluginsd(input, words, PLUGINSD_MAX_WORDS);
188
const char *command = get_word(words, num_words, 0);
189
collectors/systemd-journal.plugin/systemd-journal.c
+1
-1
@@ -690,7 +690,7 @@ static void netdata_systemd_journal_rich_message(FACETS *facets __maybe_unused,
690
static void function_systemd_journal(const char *transaction, char *function, int timeout, bool *cancelled) {
691
BUFFER *wb = buffer_create(0, NULL);
692
buffer_flush(wb);
693
- buffer_json_initialize(wb, "\"", "\"", 0, true, BUFFER_JSON_OPTIONS_NEWLINE_ON_ARRAY_ITEMS);
693
+ buffer_json_initialize(wb, "\"", "\"", 0, true, BUFFER_JSON_OPTIONS_MINIFY);
694
695
FACETS *facets = facets_create(50, FACETS_OPTION_ALL_KEYS_FTS,
696
SYSTEMD_ALWAYS_VISIBLE_KEYS,
libnetdata/buffer/buffer.c
-10
@@ -20,16 +20,6 @@ void buffer_reset(BUFFER *wb) {
20
buffer_overflow_check(wb);
21
}
22
23
-const char *buffer_tostring(BUFFER *wb)
24
-{
25
- buffer_need_bytes(wb, 1);
26
- wb->buffer[wb->len] = '\0';
27
-
28
- buffer_overflow_check(wb);
29
-
30
- return(wb->buffer);
31
-}
32
-
23
void buffer_char_replace(BUFFER *wb, char from, char to) {
24
char *s = wb->buffer, *end = &wb->buffer[wb->len];
25
libnetdata/buffer/buffer.h
+10
-1
@@ -97,7 +97,6 @@ typedef struct web_buffer {
97
#define buffer_no_cacheable(wb) do { (wb)->options |= WB_CONTENT_NO_CACHEABLE; if((wb)->options & WB_CONTENT_CACHEABLE) (wb)->options &= ~WB_CONTENT_CACHEABLE; (wb)->expires = 0; } while(0)
98
99
#define buffer_strlen(wb) ((wb)->len)
100
-const char *buffer_tostring(BUFFER *wb);
100
101
#define BUFFER_OVERFLOW_EOF "EOF"
102
@@ -158,6 +157,16 @@ void buffer_json_initialize(BUFFER *wb, const char *key_quote, const char *value
157
158
void buffer_json_finalize(BUFFER *wb);
159
160
+static const char *buffer_tostring(BUFFER *wb)
161
+{
162
+ buffer_need_bytes(wb, 1);
163
+ wb->buffer[wb->len] = '\0';
164
+
165
+ buffer_overflow_check(wb);
166
+
167
+ return(wb->buffer);
168
+}
169
+
170
static inline void _buffer_json_depth_push(BUFFER *wb, BUFFER_JSON_NODE_TYPE type) {
171
#ifdef NETDATA_INTERNAL_CHECKS
172
assert(wb->json.depth <= BUFFER_JSON_MAX_DEPTH && "BUFFER JSON: max nesting reached");
streaming/receiver.c
+22
-26
@@ -226,53 +226,47 @@ static inline bool receiver_read_compressed(struct receiver_state *r) {
226
/* Produce a full line if one exists, statefully return where we start next time.
227
* When we hit the end of the buffer with a partial line move it to the beginning for the next fill.
228
*/
229
-inline char *buffered_reader_next_line(struct buffered_reader *reader, char *dst, size_t dst_size) {
229
+inline bool buffered_reader_next_line(struct buffered_reader *reader, BUFFER *dst) {
230
+ buffer_need_bytes(dst, reader->read_len - reader->pos + 2);
231
+
232
size_t start = reader->pos;
233
234
char *ss = &reader->read_buffer[start];
235
char *se = &reader->read_buffer[reader->read_len];
234
- char *ds = dst;
235
- char *de = &dst[dst_size - 2];
236
+ char *ds = &dst->buffer[dst->len];
237
+ char *de = &ds[dst->size - dst->len - 2];
238
239
if(ss >= se) {
240
*ds = '\0';
241
reader->pos = 0;
242
reader->read_len = 0;
243
reader->read_buffer[reader->read_len] = '\0';
242
- return NULL;
244
+ return false;
245
}
246
247
// copy all bytes to buffer
246
- while(ss < se && ds < de && *ss != '\n')
248
+ while(ss < se && ds < de && *ss != '\n') {
249
*ds++ = *ss++;
250
+ dst->len++;
251
+ }
252
253
// if we have a newline, return the buffer
254
if(ss < se && ds < de && *ss == '\n') {
255
// newline found in the r->read_buffer
256
257
*ds++ = *ss++; // copy the newline too
254
- *ds = '\0';
258
+ dst->len++;
259
256
- reader->pos = ss - reader->read_buffer;
257
- return dst;
258
- }
259
-
260
- // if the destination is full, oops!
261
- if(ds == de) {
262
- netdata_log_error("STREAM: received line exceeds %d bytes. Truncating it.", PLUGINSD_LINE_MAX);
260
*ds = '\0';
261
+
262
reader->pos = ss - reader->read_buffer;
265
- return dst;
263
+ return true;
264
}
265
268
- // no newline found in the r->read_buffer
269
- // move everything to the beginning
270
- memmove(reader->read_buffer, &reader->read_buffer[start], reader->read_len - start);
271
- reader->read_len -= (int)start;
272
- reader->read_buffer[reader->read_len] = '\0';
273
- *ds = '\0';
266
reader->pos = 0;
275
- return NULL;
267
+ reader->read_len = 0;
268
+ reader->read_buffer[reader->read_len] = '\0';
269
+ return false;
270
}
271
272
bool plugin_is_enabled(struct plugind *cd);
@@ -342,10 +336,10 @@ static size_t streaming_parser(struct receiver_state *rpt, struct plugind *cd, i
336
337
buffered_reader_init(&rpt->reader);
338
345
- char buffer[PLUGINSD_LINE_MAX + 2] = "";
339
+ BUFFER *buffer = buffer_create(sizeof(rpt->reader.read_buffer), NULL);
340
while(!receiver_should_stop(rpt)) {
341
348
- if(!buffered_reader_next_line(&rpt->reader, buffer, PLUGINSD_LINE_MAX + 2)) {
342
+ if(!buffered_reader_next_line(&rpt->reader, buffer)) {
343
bool have_new_data = compressed_connection ? receiver_read_compressed(rpt) : receiver_read_uncompressed(rpt);
344
345
if(unlikely(!have_new_data)) {
@@ -356,13 +350,15 @@ static size_t streaming_parser(struct receiver_state *rpt, struct plugind *cd, i
350
continue;
351
}
352
359
- if (unlikely(parser_action(parser, buffer))) {
360
- internal_error(true, "parser_action() failed on keyword '%s'.", buffer);
353
+ if (unlikely(parser_action(parser, buffer->buffer))) {
354
receiver_set_exit_reason(rpt, STREAM_HANDSHAKE_DISCONNECT_PARSER_FAILED, false);
355
break;
356
}
364
- }
357
358
+ buffer->len = 0;
359
+ buffer->buffer[0] = '\0';
360
+ }
361
+ buffer_free(buffer);
362
result = parser->user.data_collections_count;
363
364
// free parser with the pop function
streaming/rrdpush.h
+1
-1
@@ -357,7 +357,7 @@ struct buffered_reader {
357
char read_buffer[PLUGINSD_LINE_MAX + 1];
358
};
359
360
-char *buffered_reader_next_line(struct buffered_reader *reader, char *dst, size_t dst_size);
360
+bool buffered_reader_next_line(struct buffered_reader *reader, BUFFER *dst);
361
static inline void buffered_reader_init(struct buffered_reader *reader) {
362
reader->read_buffer[0] = '\0';
363
reader->read_len = 0;