master
h 152 lines 4.69 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 #include "../libnetdata.h"
4
5 #ifndef NETDATA_BUFFERED_READER_H
6 #define NETDATA_BUFFERED_READER_H
7
8 struct buffered_reader {
9 ssize_t read_len;
10 ssize_t pos;
11 char read_buffer[PLUGINSD_LINE_MAX + 1];
12 };
13
14 static inline void buffered_reader_init(struct buffered_reader *reader) {
15 reader->read_buffer[0] = '\0';
16 reader->read_len = 0;
17 reader->pos = 0;
18 }
19
20 typedef enum {
21 BUFFERED_READER_READ_OK = 0,
22 BUFFERED_READER_READ_FAILED = -1,
23 BUFFERED_READER_READ_BUFFER_FULL = -2,
24 BUFFERED_READER_READ_POLLERR = -3,
25 BUFFERED_READER_READ_POLLHUP = -4,
26 BUFFERED_READER_READ_POLLNVAL = -5,
27 BUFFERED_READER_READ_POLL_UNKNOWN = -6,
28 BUFFERED_READER_READ_POLL_TIMEOUT = -7,
29 BUFFERED_READER_READ_POLL_CANCELLED = -8,
30 } buffered_reader_ret_t;
31
32
33 static inline buffered_reader_ret_t buffered_reader_read(struct buffered_reader *reader, int fd) {
34 #ifdef NETDATA_INTERNAL_CHECKS
35 if(reader->read_buffer[reader->read_len] != '\0')
36 fatal("read_buffer does not start with zero");
37 #endif
38
39 char *read_at = reader->read_buffer + reader->read_len;
40 ssize_t remaining = sizeof(reader->read_buffer) - reader->read_len - 1;
41
42 if(unlikely(remaining <= 0))
43 return BUFFERED_READER_READ_BUFFER_FULL;
44
45 ssize_t bytes_read = read(fd, read_at, remaining);
46 if(unlikely(bytes_read <= 0))
47 return BUFFERED_READER_READ_FAILED;
48
49 reader->read_len += bytes_read;
50 reader->read_buffer[reader->read_len] = '\0';
51
52 return BUFFERED_READER_READ_OK;
53 }
54
55 static inline buffered_reader_ret_t buffered_reader_read_timeout(struct buffered_reader *reader, int fd, int timeout_ms, bool log_error) {
56 short int revents = 0;
57 switch(wait_on_socket_or_cancel_with_timeout(
58 NULL,
59 fd, timeout_ms, POLLIN, &revents)) {
60
61 case 0: // data are waiting
62 return buffered_reader_read(reader, fd);
63
64 case 1: // timeout reached
65 if(log_error)
66 netdata_log_error("PARSER: timeout while waiting for data.");
67 return BUFFERED_READER_READ_POLL_TIMEOUT;
68
69 case -1: // thread cancelled
70 netdata_log_error("PARSER: thread cancelled while waiting for data.");
71 return BUFFERED_READER_READ_POLL_CANCELLED;
72
73 default:
74 case 2: // error on socket
75 if(revents & POLLERR) {
76 if(log_error)
77 netdata_log_error("PARSER: read failed: POLLERR.");
78 return BUFFERED_READER_READ_POLLERR;
79 }
80 if(revents & POLLHUP) {
81 if(log_error)
82 netdata_log_error("PARSER: read failed: POLLHUP.");
83 return BUFFERED_READER_READ_POLLHUP;
84 }
85 if(revents & POLLNVAL) {
86 if(log_error)
87 netdata_log_error("PARSER: read failed: POLLNVAL.");
88 return BUFFERED_READER_READ_POLLNVAL;
89 }
90 }
91
92 if(log_error)
93 netdata_log_error("PARSER: poll() returned positive number, but POLLIN|POLLERR|POLLHUP|POLLNVAL are not set.");
94 return BUFFERED_READER_READ_POLL_UNKNOWN;
95 }
96
97 /* Produce a full line if one exists, statefully return where we start next time.
98 * When we hit the end of the buffer with a partial line move it to the beginning for the next fill.
99 */
100 static inline bool buffered_reader_next_line(struct buffered_reader *reader, BUFFER *dst) {
101 buffer_need_bytes(dst, reader->read_len - reader->pos + 2);
102
103 size_t start = reader->pos;
104
105 char *ss = &reader->read_buffer[start];
106 char *se = &reader->read_buffer[reader->read_len];
107 char *ds = &dst->buffer[dst->len];
108 char *de = &ds[dst->size - dst->len - 2];
109
110 if(ss >= se) {
111 *ds = '\0';
112 reader->pos = 0;
113 reader->read_len = 0;
114 reader->read_buffer[reader->read_len] = '\0';
115 return false;
116 }
117
118 // Find out how many bytes we want to copy and whether we found a newline
119 size_t bytes_to_copy;
120 bool found_newline = false;
121 {
122 char *next_newline = (char *) memchr(ss, '\n', se - ss);
123 if (!next_newline) {
124 bytes_to_copy = se - ss;
125 } else {
126 bytes_to_copy = (next_newline - ss) + 1;
127 found_newline = true;
128 }
129 }
130
131 // Check we don't overflow the destination buffer
132 if (bytes_to_copy > (size_t)(de - ds)) {
133 bytes_to_copy = de - ds;
134 found_newline = false;
135 }
136
137 memcpy(ds, ss, bytes_to_copy);
138 ds[bytes_to_copy] = '\0';
139 dst->len += bytes_to_copy;
140
141 if (found_newline) {
142 reader->pos = start + bytes_to_copy;
143 return true;
144 }
145
146 reader->pos = 0;
147 reader->read_len = 0;
148 reader->read_buffer[reader->read_len] = '\0';
149 return false;
150 }
151
152 #endif //NETDATA_BUFFERED_READER_H