@cryptotaxi247 / netdata-1 / commits / 4f6bb2e96

switch from HTTP to stream data

Costa Tsaousis (ktsaou) committed Feb 21, 2017 at 01:47 UTC 4f6bb2e96f622233e130b249b38c99a47c3b766f
4 files changed +165 -16
src/rrdpush.c
+84 -16
@@ -147,47 +147,115 @@ void *central_netdata_push_thread(void *ptr) {
147 size_t sent_bytes = 0;
148 size_t sent_connection = 0;
149 int sock = -1;
150 - char buffer[1];
150 + char buffer[1000 + 1];
151 +
152 + struct pollfd fds[2], *ifd, *ofd;
153 +
154 + ifd = &fds[0];
155 + ofd = &fds[1];
156 +
157 + ifd->fd = rrdpush_pipe[PIPE_READ];
158 + ifd->events = POLLIN;
159 + ofd->events = POLLOUT;
160 +
161 + nfds_t fdmax = 2;
162
163 for(;;) {
164 + if(netdata_exit) break;
165 +
166 if(unlikely(sock == -1)) {
167 info("PUSH: connecting to central netdata at: %s", central_netdata_to_push_data);
168 sock = connect_to_one_of(central_netdata_to_push_data, 19999, &tv, &reconnects_counter);
169
157 - if(unlikely(sock != -1)) {
158 - info("PUSH: connected to central netdata at: %s", central_netdata_to_push_data);
170 + if(unlikely(sock == -1)) {
171 + error("PUSH: failed to connect to central netdata at: %s", central_netdata_to_push_data);
172 + sleep(5);
173 + continue;
174 + }
175 +
176 + info("PUSH: connected to central netdata at: %s", central_netdata_to_push_data);
177
160 - if(fcntl(sock, F_SETFL, O_NONBLOCK) < 0)
161 - error("PUSH: cannot set non-blocking mode for socket.");
178 + char http[1000 + 1];
179 + snprintfz(http, 1000, "GET /stream?key=%s HTTP/1.1\r\nUser-Agent: netdata-push-service/%s\r\nAccept: */*\r\n\r\n", config_get("global", "central netdata api key", ""), program_version);
180 + if(send_timeout(sock, http, strlen(http), 0, 60) == -1) {
181 + close(sock);
182 + sock = -1;
183 + error("PUSH: failed to send http header to netdata at: %s", central_netdata_to_push_data);
184 + sleep(5);
185 + continue;
186 }
163 - else
164 - error("PUSH: failed to connect to central netdata at: %s", central_netdata_to_push_data);
187 +
188 + if(recv_timeout(sock, http, 1000, 0, 60) == -1) {
189 + close(sock);
190 + sock = -1;
191 + error("PUSH: failed to receive OK from netdata at: %s", central_netdata_to_push_data);
192 + sleep(5);
193 + continue;
194 + }
195 +
196 + if(strncmp(http, "STREAM", 6)) {
197 + close(sock);
198 + sock = -1;
199 + error("PUSH: netdata servers at %s, did not send STREAM", central_netdata_to_push_data);
200 + sleep(5);
201 + continue;
202 + }
203 +
204 + if(fcntl(sock, F_SETFL, O_NONBLOCK) < 0)
205 + error("PUSH: cannot set non-blocking mode for socket.");
206
207 rrdpush_lock();
208 if(buffer_strlen(rrdpush_buffer))
209 error("PUSH: discarding %zu bytes of metrics data already in the buffer.", buffer_strlen(rrdpush_buffer));
210
211 buffer_flush(rrdpush_buffer);
171 - buffer_sprintf(rrdpush_buffer, "GET /stream?key=%s HTTP/1.1\r\nUser-Agent: netdata-push-service/%s\r\nAccept: */*\r\n\r\n", config_get("global", "central netdata api key", ""), VERSION);
212 reset_all_charts();
213 + last_host = NULL;
214 rrdpush_unlock();
215 sent_connection = 0;
216 }
217
177 - if(read(rrdpush_pipe[PIPE_READ], buffer, 1) == -1) {
178 - error("PUSH: Cannot read from internal pipe.");
179 - sleep(1);
218 + ifd->revents = 0;
219 + ofd->revents = 0;
220 + ofd->fd = sock;
221 +
222 + if(begin < buffer_strlen(rrdpush_buffer))
223 + ofd->events = POLLOUT;
224 + else
225 + ofd->events = 0;
226 +
227 + if(netdata_exit) break;
228 + int retval = poll(fds, fdmax, 60 * 1000);
229 + if(netdata_exit) break;
230 +
231 + if(unlikely(retval == -1)) {
232 + if(errno == EAGAIN || errno == EINTR)
233 + continue;
234 +
235 + error("PUSH: Failed to poll().");
236 + close(sock);
237 + sock = -1;
238 + break;
239 + }
240 + else if(unlikely(!retval)) {
241 + // timeout
242 + continue;
243 + }
244 +
245 + if(ifd->revents & POLLIN) {
246 + if(read(rrdpush_pipe[PIPE_READ], buffer, 1000) == -1)
247 + error("PUSH: Cannot read from internal pipe.");
248 }
249
182 - if(likely(sock != -1 && begin < rrdpush_buffer->len)) {
250 + if(ofd->revents & POLLOUT && begin < buffer_strlen(rrdpush_buffer)) {
251 // fprintf(stderr, "PUSH BEGIN\n");
184 - // fwrite(&rrdpush_buffer->buffer[begin], 1, rrdpush_buffer->len - begin, stderr);
252 + // fwrite(&rrdpush_buffer->buffer[begin], 1, buffer_strlen(rrdpush_buffer) - begin, stderr);
253 // fprintf(stderr, "\nPUSH END\n");
254
255 rrdpush_lock();
188 - ssize_t ret = send(sock, &rrdpush_buffer->buffer[begin], rrdpush_buffer->len - begin, MSG_DONTWAIT);
256 + ssize_t ret = send(sock, &rrdpush_buffer->buffer[begin], buffer_strlen(rrdpush_buffer) - begin, MSG_DONTWAIT);
257 if(ret == -1) {
190 - if(errno != EAGAIN) {
258 + if(errno != EAGAIN && errno != EINTR) {
259 error("PUSH: failed to send metrics to central netdata at %s. We have sent %zu bytes on this connection.", central_netdata_to_push_data, sent_connection);
260 close(sock);
261 sock = -1;
@@ -197,7 +265,7 @@ void *central_netdata_push_thread(void *ptr) {
265 sent_connection += ret;
266 sent_bytes += ret;
267 begin += ret;
200 - if(begin == rrdpush_buffer->len) {
268 + if(begin == buffer_strlen(rrdpush_buffer)) {
269 buffer_flush(rrdpush_buffer);
270 begin = 0;
271 }
src/socket.c
+62
@@ -206,3 +206,65 @@ int connect_to_one_of(const char *destination, int default_port, struct timeval
206
207 return sock;
208 }
209 +
210 +ssize_t recv_timeout(int sockfd, void *buf, size_t len, int flags, int timeout) {
211 + for(;;) {
212 + struct pollfd fd = {
213 + .fd = sockfd,
214 + .events = POLLIN,
215 + .revents = 0
216 + };
217 +
218 + errno = 0;
219 + int retval = poll(&fd, 1, timeout * 1000);
220 +
221 + if(retval == -1) {
222 + // failed
223 +
224 + if(errno == EINTR || errno == EAGAIN)
225 + continue;
226 +
227 + return -1;
228 + }
229 +
230 + if(!retval) {
231 + // timeout
232 + return 0;
233 + }
234 +
235 + if(fd.events & POLLIN) break;
236 + }
237 +
238 + return recv(sockfd, buf, len, flags);
239 +}
240 +
241 +ssize_t send_timeout(int sockfd, void *buf, size_t len, int flags, int timeout) {
242 + for(;;) {
243 + struct pollfd fd = {
244 + .fd = sockfd,
245 + .events = POLLOUT,
246 + .revents = 0
247 + };
248 +
249 + errno = 0;
250 + int retval = poll(&fd, 1, timeout * 1000);
251 +
252 + if(retval == -1) {
253 + // failed
254 +
255 + if(errno == EINTR || errno == EAGAIN)
256 + continue;
257 +
258 + return -1;
259 + }
260 +
261 + if(!retval) {
262 + // timeout
263 + return 0;
264 + }
265 +
266 + if(fd.events & POLLOUT) break;
267 + }
268 +
269 + return send(sockfd, buf, len, flags);
270 +}
src/socket.h
+3
@@ -8,4 +8,7 @@
8 extern int connect_to(const char *definition, int default_port, struct timeval *timeout);
9 extern int connect_to_one_of(const char *destination, int default_port, struct timeval *timeout, size_t *reconnects_counter);
10
11 +extern ssize_t recv_timeout(int sockfd, void *buf, size_t len, int flags, int timeout);
12 +extern ssize_t send_timeout(int sockfd, void *buf, size_t len, int flags, int timeout);
13 +
14 #endif //NETDATA_SOCKET_H
src/web_client.c
+16
@@ -1719,10 +1719,26 @@ int web_client_stream_request(RRDHOST *host, struct web_client *w, char *url) {
1719 snprintfz(cd.fullfilename, FILENAME_MAX, "%s:%s", w->client_ip, w->client_port);
1720 snprintfz(cd.cmd, PLUGINSD_CMD_MAX, "%s:%s", w->client_ip, w->client_port);
1721
1722 + if(send_timeout(w->ifd, "STREAM", 6, 0, 60) != 6) {
1723 + error("Cannot send STREAM to netdata at %s:%s", w->client_ip, w->client_port);
1724 + buffer_flush(w->response.data);
1725 + buffer_sprintf(w->response.data, "STREAM failed to reply back with STREAM");
1726 + return 400;
1727 + }
1728 +
1729 // remove the non-blocking flag from the socket
1730 if(fcntl(w->ifd, F_SETFL, fcntl(w->ifd, F_GETFL, 0) & ~O_NONBLOCK) == -1)
1731 error("STREAM from '%s:%s': cannot remove the non-blocking flag from socket %d", w->client_ip, w->client_port, w->ifd);
1732
1733 + /*
1734 + char buffer[1000 + 1];
1735 + ssize_t len;
1736 + while((len = read(w->ifd, buffer, 1000)) != -1) {
1737 + buffer[len] = '\0';
1738 + fprintf(stderr, "BEGIN READ %zu bytes\n%s\nEND READ\n", (size_t)len, buffer);
1739 + }
1740 + */
1741 +
1742 // convert the socket to a FILE *
1743 FILE *fp = fdopen(w->ifd, "r");
1744 if(!fp) {