| 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 |